mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 20:19:14 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 748bcd83fd | |||
| c93a9d0d27 |
@@ -37,6 +37,16 @@ pub mod bucket {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub mod lifecycle {
|
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 mod bucket_lifecycle_audit {
|
||||||
pub use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc;
|
pub use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -589,33 +589,192 @@ impl ExpiryOp for FreeVersionTask {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
|
enum TransitionDeleteVersionPlan {
|
||||||
|
Direct { version_id_exact: bool },
|
||||||
|
ProbeLegacyUnknown,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn legacy_transition_version_state_missing(oi: &ObjectInfo) -> Result<bool, std::io::Error> {
|
||||||
|
use rustfs_utils::http::metadata_compat::{
|
||||||
|
SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_str, get_consistent_str,
|
||||||
|
};
|
||||||
|
|
||||||
|
if !contains_key_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_STATE) {
|
||||||
|
let version_key_present = contains_key_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_ID);
|
||||||
|
if version_key_present {
|
||||||
|
if oi.transitioned_object.version_id.is_empty() {
|
||||||
|
let has_non_empty_version = oi.user_defined.iter().any(|(key, value)| {
|
||||||
|
rustfs_utils::http::metadata_compat::strip_internal_prefix_preserving_case(key)
|
||||||
|
.is_some_and(|suffix| suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID))
|
||||||
|
&& !value.is_empty()
|
||||||
|
});
|
||||||
|
if !has_non_empty_version {
|
||||||
|
// MinIO writes the transitioned-versionID key with an empty value
|
||||||
|
// for unversioned tier objects. The backend probe remains the proof.
|
||||||
|
return Ok(true);
|
||||||
|
}
|
||||||
|
} else if get_consistent_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_ID)
|
||||||
|
== Some(oi.transitioned_object.version_id.as_str())
|
||||||
|
{
|
||||||
|
return Ok(true);
|
||||||
|
}
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"legacy remote tier version metadata is conflicting or malformed",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
if !oi.transitioned_object.version_id.is_empty() {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"legacy remote tier version metadata is missing or inconsistent",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
return Ok(true);
|
||||||
|
}
|
||||||
|
let persisted = get_consistent_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_STATE).ok_or_else(|| {
|
||||||
|
std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier object has conflicting transition version state metadata",
|
||||||
|
)
|
||||||
|
})?;
|
||||||
|
if persisted != oi.transition_version_state.as_str() {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier object transition version state metadata changed during decoding",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
Ok(false)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn transition_remote_version_delete_plan(oi: &ObjectInfo) -> Result<TransitionDeleteVersionPlan, std::io::Error> {
|
||||||
|
match oi.transition_version_state {
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown => {
|
||||||
|
if legacy_transition_version_state_missing(oi)? {
|
||||||
|
Ok(TransitionDeleteVersionPlan::ProbeLegacyUnknown)
|
||||||
|
} else {
|
||||||
|
validate_transition_remote_version(oi)
|
||||||
|
.map(|version_id_exact| TransitionDeleteVersionPlan::Direct { version_id_exact })
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ => validate_transition_remote_version(oi)
|
||||||
|
.map(|version_id_exact| TransitionDeleteVersionPlan::Direct { version_id_exact }),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
|
struct ResolvedTransitionDeleteVersion {
|
||||||
|
version_id_exact: bool,
|
||||||
|
verify_missing_after_delete: bool,
|
||||||
|
remote_already_missing: bool,
|
||||||
|
}
|
||||||
|
|
||||||
async fn acquire_free_version_tier_lease(
|
async fn acquire_free_version_tier_lease(
|
||||||
oi: &ObjectInfo,
|
oi: &ObjectInfo,
|
||||||
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
||||||
) -> Result<(TierOperationLease, bool), std::io::Error> {
|
) -> Result<(TierOperationLease, TransitionDeleteVersionPlan), std::io::Error> {
|
||||||
let version_id_exact = validate_transition_remote_version(oi)?;
|
let delete_plan = transition_remote_version_delete_plan(oi)?;
|
||||||
let identity = tier_destination_id_from_metadata(&oi.user_defined)?
|
let identity = tier_destination_id_from_metadata(&oi.user_defined)?
|
||||||
.ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?;
|
.ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?;
|
||||||
let lease =
|
let lease =
|
||||||
TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, &oi.transitioned_object.tier, identity)
|
TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, &oi.transitioned_object.tier, identity)
|
||||||
.await
|
.await
|
||||||
.map_err(std::io::Error::other)?;
|
.map_err(std::io::Error::other)?;
|
||||||
Ok((lease, version_id_exact))
|
Ok((lease, delete_plan))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn resolve_transition_delete_version_plan(
|
||||||
|
oi: &ObjectInfo,
|
||||||
|
lease: &TierOperationLease,
|
||||||
|
delete_plan: TransitionDeleteVersionPlan,
|
||||||
|
) -> Result<ResolvedTransitionDeleteVersion, std::io::Error> {
|
||||||
|
match delete_plan {
|
||||||
|
TransitionDeleteVersionPlan::Direct { version_id_exact } => Ok(ResolvedTransitionDeleteVersion {
|
||||||
|
version_id_exact,
|
||||||
|
verify_missing_after_delete: false,
|
||||||
|
remote_already_missing: false,
|
||||||
|
}),
|
||||||
|
TransitionDeleteVersionPlan::ProbeLegacyUnknown => {
|
||||||
|
let probe = lease.probe_transition_candidate(&oi.transitioned_object.name).await?;
|
||||||
|
match (oi.transitioned_object.version_id.as_str(), probe) {
|
||||||
|
("", crate::services::tier::warm_backend::TransitionCandidateProbe::UnversionedPresent) => {
|
||||||
|
Ok(ResolvedTransitionDeleteVersion {
|
||||||
|
version_id_exact: false,
|
||||||
|
verify_missing_after_delete: true,
|
||||||
|
remote_already_missing: false,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
(expected, crate::services::tier::warm_backend::TransitionCandidateProbe::VersionedPresent(actual))
|
||||||
|
if !expected.is_empty() && expected == actual =>
|
||||||
|
{
|
||||||
|
lease.validate_remote_version_id(expected)?;
|
||||||
|
Ok(ResolvedTransitionDeleteVersion {
|
||||||
|
version_id_exact: true,
|
||||||
|
verify_missing_after_delete: false,
|
||||||
|
remote_already_missing: false,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
("null", crate::services::tier::warm_backend::TransitionCandidateProbe::SuspendedNullPresent) => {
|
||||||
|
lease.validate_remote_version_id("null")?;
|
||||||
|
Ok(ResolvedTransitionDeleteVersion {
|
||||||
|
version_id_exact: true,
|
||||||
|
verify_missing_after_delete: false,
|
||||||
|
remote_already_missing: false,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
(_, crate::services::tier::warm_backend::TransitionCandidateProbe::Missing) => {
|
||||||
|
Ok(ResolvedTransitionDeleteVersion {
|
||||||
|
version_id_exact: false,
|
||||||
|
verify_missing_after_delete: false,
|
||||||
|
remote_already_missing: true,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
(_, crate::services::tier::warm_backend::TransitionCandidateProbe::Unsupported) => Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::Unsupported,
|
||||||
|
"remote tier cannot prove legacy transition delete state",
|
||||||
|
)),
|
||||||
|
_ => Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::WouldBlock,
|
||||||
|
"remote tier object version state is unknown",
|
||||||
|
)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn execute_resolved_transition_delete(
|
||||||
|
oi: &ObjectInfo,
|
||||||
|
lease: &TierOperationLease,
|
||||||
|
resolved: ResolvedTransitionDeleteVersion,
|
||||||
|
) -> Result<(), std::io::Error> {
|
||||||
|
if !resolved.remote_already_missing {
|
||||||
|
delete_object_from_remote_tier_with_lease_idempotent(
|
||||||
|
&oi.transitioned_object.name,
|
||||||
|
&oi.transitioned_object.version_id,
|
||||||
|
lease,
|
||||||
|
resolved.version_id_exact,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
if resolved.verify_missing_after_delete
|
||||||
|
&& lease.probe_transition_candidate(&oi.transitioned_object.name).await?
|
||||||
|
!= crate::services::tier::warm_backend::TransitionCandidateProbe::Missing
|
||||||
|
{
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::WouldBlock,
|
||||||
|
"remote tier could not confirm legacy unversioned deletion",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_free_version_remote_object_with_lease(
|
async fn delete_free_version_remote_object_with_lease(
|
||||||
oi: &ObjectInfo,
|
oi: &ObjectInfo,
|
||||||
lease: &TierOperationLease,
|
lease: &TierOperationLease,
|
||||||
version_id_exact: bool,
|
delete_plan: TransitionDeleteVersionPlan,
|
||||||
) -> Result<(), std::io::Error> {
|
) -> Result<(), std::io::Error> {
|
||||||
delete_object_from_remote_tier_with_lease_idempotent(
|
let resolved = resolve_transition_delete_version_plan(oi, lease, delete_plan).await?;
|
||||||
&oi.transitioned_object.name,
|
execute_resolved_transition_delete(oi, lease, resolved).await
|
||||||
&oi.transitioned_object.version_id,
|
|
||||||
lease,
|
|
||||||
version_id_exact,
|
|
||||||
)
|
|
||||||
.await?;
|
|
||||||
Ok(())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn free_version_physical_topology_generation(api: &ECStore) -> String {
|
fn free_version_physical_topology_generation(api: &ECStore) -> String {
|
||||||
@@ -646,6 +805,16 @@ fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectIn
|
|||||||
if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
||||||
|| expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
|| expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
||||||
{
|
{
|
||||||
|
let candidate_legacy_missing = legacy_transition_version_state_missing(candidate)?;
|
||||||
|
let expected_legacy_missing = legacy_transition_version_state_missing(expected)?;
|
||||||
|
if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
||||||
|
&& expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
||||||
|
&& candidate_legacy_missing
|
||||||
|
&& expected_legacy_missing
|
||||||
|
&& candidate.transitioned_object.version_id == expected.transitioned_object.version_id
|
||||||
|
{
|
||||||
|
return Ok(true);
|
||||||
|
}
|
||||||
return Err(std::io::Error::new(
|
return Err(std::io::Error::new(
|
||||||
std::io::ErrorKind::WouldBlock,
|
std::io::ErrorKind::WouldBlock,
|
||||||
"tier free-version remote version state is unknown",
|
"tier free-version remote version state is unknown",
|
||||||
@@ -721,7 +890,7 @@ async fn cleanup_free_version_exact(api: Arc<ECStore>, oi: &ObjectInfo, cancel:
|
|||||||
.acquire_bucket_lifecycle_read_lock(&oi.bucket)
|
.acquire_bucket_lifecycle_read_lock(&oi.bucket)
|
||||||
.await
|
.await
|
||||||
.map_err(std::io::Error::other)?;
|
.map_err(std::io::Error::other)?;
|
||||||
let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?;
|
let (lease, delete_plan) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?;
|
||||||
let local_object = encode_dir_object(&oi.name);
|
let local_object = encode_dir_object(&oi.name);
|
||||||
let object_guards = api
|
let object_guards = api
|
||||||
.acquire_all_physical_object_write_locks("tier_free_version_cleanup", &oi.bucket, &local_object)
|
.acquire_all_physical_object_write_locks("tier_free_version_cleanup", &oi.bucket, &local_object)
|
||||||
@@ -739,16 +908,30 @@ async fn cleanup_free_version_exact(api: Arc<ECStore>, oi: &ObjectInfo, cancel:
|
|||||||
"tier free-version cleanup fence is invalid before remote delete",
|
"tier free-version cleanup fence is invalid before remote delete",
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
let resolved = tokio::select! {
|
||||||
|
_ = cancel.cancelled() => {
|
||||||
|
return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled"));
|
||||||
|
}
|
||||||
|
result = tokio::time::timeout_at(deadline, resolve_transition_delete_version_plan(oi, &lease, delete_plan)) => {
|
||||||
|
result.map_err(|_| {
|
||||||
|
std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote probe timed out")
|
||||||
|
})??
|
||||||
|
}
|
||||||
|
};
|
||||||
|
if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::WouldBlock,
|
||||||
|
"tier free-version cleanup fence changed after remote probe",
|
||||||
|
));
|
||||||
|
}
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
_ = cancel.cancelled() => {
|
_ = cancel.cancelled() => {
|
||||||
return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled"));
|
return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled"));
|
||||||
}
|
}
|
||||||
result = tokio::time::timeout_at(
|
result = tokio::time::timeout_at(deadline, execute_resolved_transition_delete(oi, &lease, resolved)) => {
|
||||||
deadline,
|
result.map_err(|_| {
|
||||||
delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact),
|
std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote delete timed out")
|
||||||
) => {
|
})??;
|
||||||
result
|
|
||||||
.map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote delete timed out"))??;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) {
|
if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) {
|
||||||
@@ -796,8 +979,8 @@ async fn delete_free_version_remote_object(
|
|||||||
oi: &ObjectInfo,
|
oi: &ObjectInfo,
|
||||||
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
||||||
) -> Result<(), std::io::Error> {
|
) -> Result<(), std::io::Error> {
|
||||||
let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?;
|
let (lease, delete_plan) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?;
|
||||||
delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact).await
|
delete_free_version_remote_object_with_lease(oi, &lease, delete_plan).await
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(
|
#[allow(
|
||||||
@@ -813,8 +996,8 @@ where
|
|||||||
F: FnOnce() -> Fut,
|
F: FnOnce() -> Fut,
|
||||||
Fut: std::future::Future<Output = T>,
|
Fut: std::future::Future<Output = T>,
|
||||||
{
|
{
|
||||||
let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?;
|
let (lease, delete_plan) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?;
|
||||||
delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact).await?;
|
delete_free_version_remote_object_with_lease(oi, &lease, delete_plan).await?;
|
||||||
let result = delete_local().await;
|
let result = delete_local().await;
|
||||||
drop(lease);
|
drop(lease);
|
||||||
Ok(result)
|
Ok(result)
|
||||||
@@ -4693,6 +4876,39 @@ fn validate_transition_remote_version(oi: &ObjectInfo) -> Result<bool, std::io::
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
|
enum TransitionReadVersionPlan {
|
||||||
|
Direct,
|
||||||
|
ProbeLegacyUnversioned,
|
||||||
|
}
|
||||||
|
|
||||||
|
const LEGACY_TRANSITION_READ_PROBE_TIMEOUT: StdDuration = StdDuration::from_secs(30);
|
||||||
|
|
||||||
|
fn transition_remote_version_read_plan(oi: &ObjectInfo) -> Result<TransitionReadVersionPlan, std::io::Error> {
|
||||||
|
let version = oi.transitioned_object.version_id.as_str();
|
||||||
|
match oi.transition_version_state {
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown => {
|
||||||
|
if !legacy_transition_version_state_missing(oi)? {
|
||||||
|
return validate_transition_remote_version(oi).map(|_| TransitionReadVersionPlan::Direct);
|
||||||
|
}
|
||||||
|
if version.is_empty() {
|
||||||
|
Ok(TransitionReadVersionPlan::ProbeLegacyUnversioned)
|
||||||
|
} else {
|
||||||
|
Ok(TransitionReadVersionPlan::Direct)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
rustfs_filemeta::TransitionVersionState::KnownDisabled if version.is_empty() => Ok(TransitionReadVersionPlan::Direct),
|
||||||
|
rustfs_filemeta::TransitionVersionState::SuspendedNull if version == "null" => Ok(TransitionReadVersionPlan::Direct),
|
||||||
|
rustfs_filemeta::TransitionVersionState::Exact if !version.is_empty() && version != "null" => {
|
||||||
|
Ok(TransitionReadVersionPlan::Direct)
|
||||||
|
}
|
||||||
|
_ => Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier object version state conflicts with its version ID",
|
||||||
|
)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// The resolver joins the tier manager as the second injected port this read
|
// The resolver joins the tier manager as the second injected port this read
|
||||||
// needs; grouping the request half into a struct would churn every call site of
|
// needs; grouping the request half into a struct would churn every call site of
|
||||||
// a bug fix.
|
// a bug fix.
|
||||||
@@ -4707,7 +4923,12 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
|||||||
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
||||||
resolver: Option<&dyn ObjectEncryptionResolver>,
|
resolver: Option<&dyn ObjectEncryptionResolver>,
|
||||||
) -> Result<GetObjectReader, std::io::Error> {
|
) -> Result<GetObjectReader, std::io::Error> {
|
||||||
validate_transition_remote_version(oi)?;
|
let read_plan = transition_remote_version_read_plan(oi)?;
|
||||||
|
// Reject invalid ranges and encryption requests before a compatibility
|
||||||
|
// probe can amplify them into remote listing work.
|
||||||
|
let plan = ReadPlan::build_for_request(rs.clone(), oi, opts, h, resolver)
|
||||||
|
.await
|
||||||
|
.map_err(|err| std::io::Error::other(format!("building the read plan for {bucket}/{object} failed: {err}")))?;
|
||||||
let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?;
|
let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?;
|
||||||
let lease = match expected_identity {
|
let lease = match expected_identity {
|
||||||
Some(identity) => {
|
Some(identity) => {
|
||||||
@@ -4721,7 +4942,36 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
|||||||
Err(err) => return Err(std::io::Error::other(err)),
|
Err(err) => return Err(std::io::Error::other(err)),
|
||||||
};
|
};
|
||||||
|
|
||||||
tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?;
|
match read_plan {
|
||||||
|
TransitionReadVersionPlan::Direct => {
|
||||||
|
tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?;
|
||||||
|
}
|
||||||
|
TransitionReadVersionPlan::ProbeLegacyUnversioned => {
|
||||||
|
// RUSTFS_COMPAT_TODO(backlog#2203): remove operation-time probing
|
||||||
|
// after an admin reconcile can persist every proven legacy state.
|
||||||
|
let probe = tokio::time::timeout(
|
||||||
|
LEGACY_TRANSITION_READ_PROBE_TIMEOUT,
|
||||||
|
tgt_client.probe_transition_candidate(&oi.transitioned_object.name),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "legacy remote tier version probe timed out"))??;
|
||||||
|
match probe {
|
||||||
|
crate::services::tier::warm_backend::TransitionCandidateProbe::UnversionedPresent => {}
|
||||||
|
crate::services::tier::warm_backend::TransitionCandidateProbe::Unsupported => {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::Unsupported,
|
||||||
|
"remote tier cannot prove legacy unversioned transition state",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
_ => {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier object version state is unknown",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// The same read plan the local path uses, so the tier fetch is positioned in
|
// The same read plan the local path uses, so the tier fetch is positioned in
|
||||||
// the object's *stored* coordinate system and the stream is handed the same
|
// the object's *stored* coordinate system and the stream is handed the same
|
||||||
@@ -4729,9 +4979,6 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
|||||||
// through a plaintext-coordinate range and skipping the transform is how a
|
// through a plaintext-coordinate range and skipping the transform is how a
|
||||||
// transitioned SSE object used to come back as silently corrupt bytes of the
|
// transitioned SSE object used to come back as silently corrupt bytes of the
|
||||||
// right length (rustfs/rustfs#6025).
|
// right length (rustfs/rustfs#6025).
|
||||||
let plan = ReadPlan::build_for_request(rs.clone(), oi, opts, h, resolver)
|
|
||||||
.await
|
|
||||||
.map_err(|err| std::io::Error::other(format!("building the read plan for {bucket}/{object} failed: {err}")))?;
|
|
||||||
let (off, length) = (plan.storage_offset() as i64, plan.storage_length());
|
let (off, length) = (plan.storage_offset() as i64, plan.storage_length());
|
||||||
let mut gopts = WarmBackendGetOpts::default();
|
let mut gopts = WarmBackendGetOpts::default();
|
||||||
|
|
||||||
@@ -5629,11 +5876,13 @@ mod tests {
|
|||||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||||
use crate::object_api::{ObjectInfo, ObjectOptions, PutObjReader};
|
use crate::object_api::{ObjectInfo, ObjectOptions, PutObjReader};
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
|
use crate::services::tier::test_util::MockWarmOp;
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
use crate::services::tier::test_util::register_mock_tier;
|
use crate::services::tier::test_util::register_mock_tier;
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
use crate::services::tier::tier::TierConfigMgr;
|
use crate::services::tier::tier::TierConfigMgr;
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
use crate::services::tier::warm_backend::WarmBackend as _;
|
use crate::services::tier::warm_backend::{TransitionCandidateProbe, WarmBackend as _};
|
||||||
use crate::set_disk::{MultipartCommitBarrier, MultipartCommitPause};
|
use crate::set_disk::{MultipartCommitBarrier, MultipartCommitPause};
|
||||||
use crate::set_disk::{RUSTFS_MULTIPART_BUCKET_KEY, RUSTFS_MULTIPART_OBJECT_KEY};
|
use crate::set_disk::{RUSTFS_MULTIPART_BUCKET_KEY, RUSTFS_MULTIPART_OBJECT_KEY};
|
||||||
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
||||||
@@ -6329,7 +6578,75 @@ mod tests {
|
|||||||
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn transitioned_get_rejects_unknown_version_state_before_backend_io() {
|
async fn transitioned_get_allows_legacy_unknown_exact_version_for_non_destructive_read() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let backend = register_mock_tier(&manager, &tier).await;
|
||||||
|
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||||
|
let body = Bytes::from_static(b"legacy transitioned object body");
|
||||||
|
let remote_version = backend
|
||||||
|
.put(
|
||||||
|
&remote_object,
|
||||||
|
ReaderImpl::Body(body.clone()),
|
||||||
|
i64::try_from(body.len()).expect("body length should fit"),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("mock remote object should be stored");
|
||||||
|
let mut user_defined = HashMap::new();
|
||||||
|
insert_legacy_transition_version_id(&mut user_defined, &remote_version);
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
size: i64::try_from(body.len()).expect("body length should fit"),
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: remote_object,
|
||||||
|
version_id: remote_version,
|
||||||
|
status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
|
tier: tier.clone(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let range = Some(crate::storage_api_contracts::range::HTTPRangeSpec {
|
||||||
|
is_suffix_length: false,
|
||||||
|
start: 7,
|
||||||
|
end: 18,
|
||||||
|
});
|
||||||
|
let mut reader = get_transitioned_object_reader_with_tier_manager(
|
||||||
|
&object_info.bucket,
|
||||||
|
&object_info.name,
|
||||||
|
&range,
|
||||||
|
&HeaderMap::new(),
|
||||||
|
&object_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&manager,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("legacy unknown state should still allow a non-destructive read");
|
||||||
|
let mut got = Vec::new();
|
||||||
|
reader
|
||||||
|
.stream
|
||||||
|
.read_to_end(&mut got)
|
||||||
|
.await
|
||||||
|
.expect("transitioned reader should drain");
|
||||||
|
|
||||||
|
assert_eq!(got, &body.as_ref()[7..=18]);
|
||||||
|
assert_eq!(backend.get_count().await, 1);
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
assert_eq!(
|
||||||
|
TierConfigMgr::active_operation_lease_count(&manager, &tier).await,
|
||||||
|
0,
|
||||||
|
"tier generation lease should release after EOF"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn transitioned_get_rejects_explicit_unknown_version_state_before_backend_io() {
|
||||||
let manager = TierConfigMgr::new();
|
let manager = TierConfigMgr::new();
|
||||||
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
let backend = register_mock_tier(&manager, &tier).await;
|
let backend = register_mock_tier(&manager, &tier).await;
|
||||||
@@ -6345,6 +6662,181 @@ mod tests {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined_with_transition_version_state(rustfs_filemeta::TransitionVersionState::Unknown).into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = match get_transitioned_object_reader_with_tier_manager(
|
||||||
|
&object_info.bucket,
|
||||||
|
&object_info.name,
|
||||||
|
&None,
|
||||||
|
&HeaderMap::new(),
|
||||||
|
&object_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&manager,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(_) => panic!("explicit unknown remote version state must fail before backend IO"),
|
||||||
|
Err(err) => err,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(backend.op_log().await, Vec::<MockWarmOp>::new());
|
||||||
|
assert_eq!(backend.get_count().await, 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn transitioned_get_rejects_present_but_invalid_legacy_version_metadata() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let backend = register_mock_tier(&manager, &tier).await;
|
||||||
|
|
||||||
|
for persisted_version in [
|
||||||
|
Uuid::nil().to_string(),
|
||||||
|
"\u{fffd}".to_string(),
|
||||||
|
"bad\u{0001}version".to_string(),
|
||||||
|
] {
|
||||||
|
let mut user_defined = HashMap::new();
|
||||||
|
insert_legacy_transition_version_id(&mut user_defined, &persisted_version);
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
size: 1,
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: "remote/object".to_string(),
|
||||||
|
version_id: String::new(),
|
||||||
|
status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
|
tier: tier.clone(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = match get_transitioned_object_reader_with_tier_manager(
|
||||||
|
&object_info.bucket,
|
||||||
|
&object_info.name,
|
||||||
|
&None,
|
||||||
|
&HeaderMap::new(),
|
||||||
|
&object_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&manager,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(_) => panic!("present but invalid legacy version metadata must fail before backend IO"),
|
||||||
|
Err(err) => err,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
}
|
||||||
|
|
||||||
|
assert_eq!(backend.op_log().await, Vec::<MockWarmOp>::new());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn transitioned_get_probes_legacy_empty_unknown_state_before_unversioned_read() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let backend = register_mock_tier(&manager, &tier).await;
|
||||||
|
backend.set_put_remote_version(Some(String::new())).await;
|
||||||
|
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||||
|
let body = Bytes::from_static(b"legacy unversioned transitioned object body");
|
||||||
|
let remote_version = backend
|
||||||
|
.put(
|
||||||
|
&remote_object,
|
||||||
|
ReaderImpl::Body(body.clone()),
|
||||||
|
i64::try_from(body.len()).expect("body length should fit"),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("mock remote object should be stored");
|
||||||
|
assert!(remote_version.is_empty());
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
size: i64::try_from(body.len()).expect("body length should fit"),
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: remote_object.clone(),
|
||||||
|
version_id: String::new(),
|
||||||
|
status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
|
tier: tier.clone(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: HashMap::from([("x-minio-internal-transitioned-versionID".to_string(), String::new())]).into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let mut reader = get_transitioned_object_reader_with_tier_manager(
|
||||||
|
&object_info.bucket,
|
||||||
|
&object_info.name,
|
||||||
|
&None,
|
||||||
|
&HeaderMap::new(),
|
||||||
|
&object_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&manager,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("probe-proven legacy unversioned state should allow a non-destructive read");
|
||||||
|
let mut got = Vec::new();
|
||||||
|
reader
|
||||||
|
.stream
|
||||||
|
.read_to_end(&mut got)
|
||||||
|
.await
|
||||||
|
.expect("transitioned reader should drain");
|
||||||
|
|
||||||
|
assert_eq!(got, body.as_ref());
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
assert_eq!(
|
||||||
|
backend.op_log().await,
|
||||||
|
vec![
|
||||||
|
MockWarmOp::Put {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
MockWarmOp::Get { object: remote_object },
|
||||||
|
]
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
TierConfigMgr::active_operation_lease_count(&manager, &tier).await,
|
||||||
|
0,
|
||||||
|
"tier generation lease should release after EOF"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn transitioned_get_rejects_ambiguous_empty_unknown_state_without_backend_get() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let backend = register_mock_tier(&manager, &tier).await;
|
||||||
|
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||||
|
backend
|
||||||
|
.set_transition_candidate_probe_override(Some(TransitionCandidateProbe::VersionedPresent(
|
||||||
|
"versioned-candidate".to_string(),
|
||||||
|
)))
|
||||||
|
.await;
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
size: 1,
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: remote_object.clone(),
|
||||||
|
version_id: String::new(),
|
||||||
|
status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
|
tier,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -6360,19 +6852,28 @@ mod tests {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(_) => panic!("unknown remote version state must fail before backend IO"),
|
Ok(_) => panic!("versioned legacy unknown state without stored version must fail before backend GET"),
|
||||||
Err(err) => err,
|
Err(err) => err,
|
||||||
};
|
};
|
||||||
|
|
||||||
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(backend.op_log().await, vec![MockWarmOp::Probe { object: remote_object }]);
|
||||||
assert_eq!(backend.get_count().await, 0);
|
assert_eq!(backend.get_count().await, 0);
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn free_version_delete_rejects_unknown_version_state_before_backend_io() {
|
async fn free_version_delete_rejects_explicit_unknown_before_backend_io() {
|
||||||
let manager = TierConfigMgr::new();
|
let manager = TierConfigMgr::new();
|
||||||
let backend = register_mock_tier(&manager, "WARM").await;
|
let backend = register_mock_tier(&manager, "WARM").await;
|
||||||
|
let identity = test_tier_destination_identity(&manager, "WARM").await;
|
||||||
|
let mut user_defined = user_defined_with_tier_destination_identity(identity);
|
||||||
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
|
&mut user_defined,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown.as_str().to_string(),
|
||||||
|
);
|
||||||
let object_info = ObjectInfo {
|
let object_info = ObjectInfo {
|
||||||
transitioned_object: TransitionedObject {
|
transitioned_object: TransitionedObject {
|
||||||
name: "remote/object".to_string(),
|
name: "remote/object".to_string(),
|
||||||
@@ -6381,14 +6882,243 @@ mod tests {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
let err = super::delete_free_version_remote_object(&object_info, &manager)
|
let err = super::delete_free_version_remote_object(&object_info, &manager)
|
||||||
.await
|
.await
|
||||||
.expect_err("unknown remote version state must fail before backend IO");
|
.expect_err("explicit unknown cleanup must fail before backend IO");
|
||||||
|
|
||||||
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert!(err.to_string().contains("version state is unknown"));
|
||||||
|
assert_eq!(backend.op_log().await, Vec::<MockWarmOp>::new());
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
async fn test_tier_destination_identity(
|
||||||
|
manager: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
|
||||||
|
tier: &str,
|
||||||
|
) -> crate::services::tier::tier::TierDestinationId {
|
||||||
|
TierConfigMgr::acquire_operation_lease(manager, tier)
|
||||||
|
.await
|
||||||
|
.expect("test tier lease should be available")
|
||||||
|
.backend_identity()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
fn user_defined_with_tier_destination_identity(
|
||||||
|
identity: crate::services::tier::tier::TierDestinationId,
|
||||||
|
) -> HashMap<String, String> {
|
||||||
|
let mut user_defined = HashMap::new();
|
||||||
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
|
&mut user_defined,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||||
|
rustfs_utils::crypto::hex(identity),
|
||||||
|
);
|
||||||
|
user_defined
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
fn user_defined_with_transition_version_state(state: rustfs_filemeta::TransitionVersionState) -> HashMap<String, String> {
|
||||||
|
let mut user_defined = HashMap::new();
|
||||||
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
|
&mut user_defined,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
state.as_str().to_string(),
|
||||||
|
);
|
||||||
|
user_defined
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
fn insert_legacy_transition_version_id(user_defined: &mut HashMap<String, String>, version_id: &str) {
|
||||||
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
|
user_defined,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_ID,
|
||||||
|
version_id.to_string(),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn free_version_tuple_rejects_mixed_legacy_missing_and_explicit_unknown() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
register_mock_tier(&manager, "WARM").await;
|
||||||
|
let identity = test_tier_destination_identity(&manager, "WARM").await;
|
||||||
|
let mut legacy_metadata = user_defined_with_tier_destination_identity(identity);
|
||||||
|
insert_legacy_transition_version_id(&mut legacy_metadata, "legacy-version");
|
||||||
|
let mut explicit_metadata = legacy_metadata.clone();
|
||||||
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
|
&mut explicit_metadata,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown.as_str().to_string(),
|
||||||
|
);
|
||||||
|
let make_info = |user_defined: HashMap<String, String>| ObjectInfo {
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: "remote/object".to_string(),
|
||||||
|
version_id: "legacy-version".to_string(),
|
||||||
|
tier: "WARM".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = super::free_version_remote_tuple_matches(&make_info(legacy_metadata), &make_info(explicit_metadata))
|
||||||
|
.expect_err("mixed legacy-missing and explicit unknown provenance must fail closed");
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn free_version_delete_probes_legacy_unknown_exact_version_before_remove() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = "WARM";
|
||||||
|
let backend = register_mock_tier(&manager, tier).await;
|
||||||
|
let identity = test_tier_destination_identity(&manager, tier).await;
|
||||||
|
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||||
|
let body = Bytes::from_static(b"legacy exact cleanup body");
|
||||||
|
let remote_version = backend
|
||||||
|
.put(
|
||||||
|
&remote_object,
|
||||||
|
ReaderImpl::Body(body),
|
||||||
|
i64::try_from(b"legacy exact cleanup body".len()).expect("body length should fit"),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("mock remote object should be stored");
|
||||||
|
let mut user_defined = user_defined_with_tier_destination_identity(identity);
|
||||||
|
insert_legacy_transition_version_id(&mut user_defined, &remote_version);
|
||||||
|
backend.clear_op_log().await;
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: remote_object.clone(),
|
||||||
|
version_id: remote_version,
|
||||||
|
tier: tier.to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
super::delete_free_version_remote_object(&object_info, &manager)
|
||||||
|
.await
|
||||||
|
.expect("probe-proven legacy exact cleanup should delete the remote version");
|
||||||
|
super::delete_free_version_remote_object(&object_info, &manager)
|
||||||
|
.await
|
||||||
|
.expect("a retry after the exact remote version is already missing should be idempotent");
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
backend.op_log().await,
|
||||||
|
vec![
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
MockWarmOp::Remove {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
]
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
backend.remove_versions().await,
|
||||||
|
vec![(remote_object, object_info.transitioned_object.version_id)]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn free_version_delete_probes_legacy_unknown_unversioned_before_remove() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = "WARM";
|
||||||
|
let backend = register_mock_tier(&manager, tier).await;
|
||||||
|
backend.set_put_remote_version(Some(String::new())).await;
|
||||||
|
let identity = test_tier_destination_identity(&manager, tier).await;
|
||||||
|
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||||
|
let body = Bytes::from_static(b"legacy unversioned cleanup body");
|
||||||
|
let remote_version = backend
|
||||||
|
.put(
|
||||||
|
&remote_object,
|
||||||
|
ReaderImpl::Body(body),
|
||||||
|
i64::try_from(b"legacy unversioned cleanup body".len()).expect("body length should fit"),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("mock remote object should be stored");
|
||||||
|
assert!(remote_version.is_empty());
|
||||||
|
backend.clear_op_log().await;
|
||||||
|
let mut user_defined = user_defined_with_tier_destination_identity(identity);
|
||||||
|
user_defined.insert("x-minio-internal-transitioned-versionID".to_string(), String::new());
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: remote_object.clone(),
|
||||||
|
version_id: String::new(),
|
||||||
|
tier: tier.to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
super::delete_free_version_remote_object(&object_info, &manager)
|
||||||
|
.await
|
||||||
|
.expect("probe-proven legacy unversioned cleanup should delete the remote object");
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
backend.op_log().await,
|
||||||
|
vec![
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
MockWarmOp::Remove {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: remote_object.clone()
|
||||||
|
},
|
||||||
|
]
|
||||||
|
);
|
||||||
|
assert_eq!(backend.remove_versions().await, vec![(remote_object, String::new())]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn free_version_delete_retains_legacy_unknown_when_probe_disagrees() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = "WARM";
|
||||||
|
let backend = register_mock_tier(&manager, tier).await;
|
||||||
|
let identity = test_tier_destination_identity(&manager, tier).await;
|
||||||
|
let remote_object = format!("remote/{}", Uuid::new_v4());
|
||||||
|
backend
|
||||||
|
.set_transition_candidate_probe_override(Some(TransitionCandidateProbe::VersionedPresent(
|
||||||
|
"different-version".to_string(),
|
||||||
|
)))
|
||||||
|
.await;
|
||||||
|
let mut user_defined = user_defined_with_tier_destination_identity(identity);
|
||||||
|
insert_legacy_transition_version_id(&mut user_defined, "legacy-version");
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: remote_object.clone(),
|
||||||
|
version_id: "legacy-version".to_string(),
|
||||||
|
tier: tier.to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
user_defined: user_defined.into(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = super::delete_free_version_remote_object(&object_info, &manager)
|
||||||
|
.await
|
||||||
|
.expect_err("legacy unknown cleanup must not delete when the probe disagrees");
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock);
|
||||||
|
assert_eq!(backend.op_log().await, vec![MockWarmOp::Probe { object: remote_object }]);
|
||||||
assert_eq!(backend.remove_count().await, 0);
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
@@ -18,6 +18,7 @@ mod config_boundary;
|
|||||||
pub mod core;
|
pub mod core;
|
||||||
mod durable_namespace;
|
mod durable_namespace;
|
||||||
pub mod evaluator;
|
pub mod evaluator;
|
||||||
|
pub mod legacy_transition_state_reconcile;
|
||||||
pub mod manual_transition_job;
|
pub mod manual_transition_job;
|
||||||
mod metadata_boundary;
|
mod metadata_boundary;
|
||||||
pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs, get_lifecycle_config};
|
pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs, get_lifecycle_config};
|
||||||
|
|||||||
@@ -830,6 +830,7 @@ impl From<TransitionCandidateProbe> for TransitionOperatorProbe {
|
|||||||
match value {
|
match value {
|
||||||
TransitionCandidateProbe::Missing => Self::Missing,
|
TransitionCandidateProbe::Missing => Self::Missing,
|
||||||
TransitionCandidateProbe::UnversionedPresent => Self::UnversionedPresent,
|
TransitionCandidateProbe::UnversionedPresent => Self::UnversionedPresent,
|
||||||
|
TransitionCandidateProbe::SuspendedNullPresent => Self::VersionedPresent("null".to_string()),
|
||||||
TransitionCandidateProbe::VersionedPresent(version_id) => Self::VersionedPresent(version_id),
|
TransitionCandidateProbe::VersionedPresent(version_id) => Self::VersionedPresent(version_id),
|
||||||
TransitionCandidateProbe::Ambiguous => Self::Ambiguous,
|
TransitionCandidateProbe::Ambiguous => Self::Ambiguous,
|
||||||
TransitionCandidateProbe::Unsupported => Self::Unsupported,
|
TransitionCandidateProbe::Unsupported => Self::Unsupported,
|
||||||
@@ -1183,6 +1184,10 @@ async fn recover_unknown_upload_outcome(
|
|||||||
TransitionCandidateProbe::UnversionedPresent => {
|
TransitionCandidateProbe::UnversionedPresent => {
|
||||||
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
|
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
|
||||||
}
|
}
|
||||||
|
TransitionCandidateProbe::SuspendedNullPresent => {
|
||||||
|
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned("null".to_string()))
|
||||||
|
.await
|
||||||
|
}
|
||||||
TransitionCandidateProbe::VersionedPresent(version_id)
|
TransitionCandidateProbe::VersionedPresent(version_id)
|
||||||
if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) =>
|
if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) =>
|
||||||
{
|
{
|
||||||
@@ -1512,6 +1517,14 @@ mod tests {
|
|||||||
|
|
||||||
const BACKEND_FINGERPRINT: [u8; 32] = [7; 32];
|
const BACKEND_FINGERPRINT: [u8; 32] = [7; 32];
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn suspended_null_probe_preserves_operator_null_version_semantics() {
|
||||||
|
assert_eq!(
|
||||||
|
TransitionOperatorProbe::from(TransitionCandidateProbe::SuspendedNullPresent),
|
||||||
|
TransitionOperatorProbe::VersionedPresent("null".to_string())
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Default)]
|
#[derive(Default)]
|
||||||
struct MemoryTransactionStore {
|
struct MemoryTransactionStore {
|
||||||
records: HashMap<Uuid, Vec<u8>>,
|
records: HashMap<Uuid, Vec<u8>>,
|
||||||
|
|||||||
@@ -701,7 +701,7 @@ impl WarmBackend for MockWarmBackend {
|
|||||||
Ok(version)
|
Ok(version)
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn get(&self, object: &str, _rv: &str, opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
|
async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
|
||||||
self.precondition().await?;
|
self.precondition().await?;
|
||||||
let barrier = self.inner.get_barrier.lock().await.take();
|
let barrier = self.inner.get_barrier.lock().await.take();
|
||||||
if let Some(barrier) = barrier {
|
if let Some(barrier) = barrier {
|
||||||
@@ -719,6 +719,9 @@ impl WarmBackend for MockWarmBackend {
|
|||||||
let Some(stored) = objects.get(object) else {
|
let Some(stored) = objects.get(object) else {
|
||||||
return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "mock object not found"));
|
return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "mock object not found"));
|
||||||
};
|
};
|
||||||
|
if !rv.is_empty() && stored.remote_version_id != rv {
|
||||||
|
return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "NoSuchVersion"));
|
||||||
|
}
|
||||||
let bytes = &stored.bytes;
|
let bytes = &stored.bytes;
|
||||||
|
|
||||||
let start = opts.start_offset.max(0) as usize;
|
let start = opts.start_offset.max(0) as usize;
|
||||||
@@ -789,6 +792,8 @@ impl WarmBackend for MockWarmBackend {
|
|||||||
};
|
};
|
||||||
if stored.remote_version_id.is_empty() {
|
if stored.remote_version_id.is_empty() {
|
||||||
Ok(TransitionCandidateProbe::UnversionedPresent)
|
Ok(TransitionCandidateProbe::UnversionedPresent)
|
||||||
|
} else if stored.remote_version_id == "null" {
|
||||||
|
Ok(TransitionCandidateProbe::SuspendedNullPresent)
|
||||||
} else {
|
} else {
|
||||||
Ok(TransitionCandidateProbe::VersionedPresent(stored.remote_version_id.clone()))
|
Ok(TransitionCandidateProbe::VersionedPresent(stored.remote_version_id.clone()))
|
||||||
}
|
}
|
||||||
@@ -988,7 +993,7 @@ mod tests {
|
|||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn mock_probe_distinguishes_missing_unversioned_and_versioned_candidates() {
|
async fn mock_probe_distinguishes_all_known_remote_version_states() {
|
||||||
let backend = MockWarmBackend::new();
|
let backend = MockWarmBackend::new();
|
||||||
|
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
@@ -1026,6 +1031,19 @@ mod tests {
|
|||||||
TransitionCandidateProbe::VersionedPresent(remote_version)
|
TransitionCandidateProbe::VersionedPresent(remote_version)
|
||||||
);
|
);
|
||||||
|
|
||||||
|
backend.set_put_remote_version(Some("null".to_string())).await;
|
||||||
|
backend
|
||||||
|
.put("suspended-null", ReaderImpl::Body(Bytes::new()), 0)
|
||||||
|
.await
|
||||||
|
.expect("put suspended null candidate");
|
||||||
|
assert_eq!(
|
||||||
|
backend
|
||||||
|
.probe_transition_candidate_state("suspended-null")
|
||||||
|
.await
|
||||||
|
.expect("probe suspended null candidate"),
|
||||||
|
TransitionCandidateProbe::SuspendedNullPresent
|
||||||
|
);
|
||||||
|
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
backend
|
backend
|
||||||
.op_log()
|
.op_log()
|
||||||
@@ -1033,7 +1051,7 @@ mod tests {
|
|||||||
.into_iter()
|
.into_iter()
|
||||||
.filter(|op| matches!(op, MockWarmOp::Probe { .. }))
|
.filter(|op| matches!(op, MockWarmOp::Probe { .. }))
|
||||||
.count(),
|
.count(),
|
||||||
3
|
4
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -78,6 +78,8 @@ pub struct WarmBackendGetOpts {
|
|||||||
pub enum TransitionCandidateProbe {
|
pub enum TransitionCandidateProbe {
|
||||||
Missing,
|
Missing,
|
||||||
UnversionedPresent,
|
UnversionedPresent,
|
||||||
|
/// The provider listed the candidate under the exact S3 `null` version.
|
||||||
|
SuspendedNullPresent,
|
||||||
VersionedPresent(String),
|
VersionedPresent(String),
|
||||||
Ambiguous,
|
Ambiguous,
|
||||||
Unsupported,
|
Unsupported,
|
||||||
|
|||||||
@@ -214,7 +214,8 @@ impl WarmBackendS3 {
|
|||||||
return Ok(TransitionCandidateProbe::Ambiguous);
|
return Ok(TransitionCandidateProbe::Ambiguous);
|
||||||
}
|
}
|
||||||
if !versions.is_truncated {
|
if !versions.is_truncated {
|
||||||
return classify_transition_candidates(candidates, bucket_versioning);
|
let confirmed_versioning = self.remote_bucket_versioning().await?;
|
||||||
|
return classify_transition_candidates(candidates, bucket_versioning, confirmed_versioning);
|
||||||
}
|
}
|
||||||
|
|
||||||
advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?;
|
advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?;
|
||||||
@@ -275,7 +276,8 @@ impl WarmBackendS3 {
|
|||||||
version_id: matched_version,
|
version_id: matched_version,
|
||||||
ambiguous: false,
|
ambiguous: false,
|
||||||
};
|
};
|
||||||
return classify_transition_candidates(candidates, bucket_versioning);
|
let confirmed_versioning = self.remote_bucket_versioning().await?;
|
||||||
|
return classify_transition_candidates(candidates, bucket_versioning, confirmed_versioning);
|
||||||
}
|
}
|
||||||
advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?;
|
advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?;
|
||||||
}
|
}
|
||||||
@@ -310,9 +312,15 @@ fn transition_candidate_metadata_matches(
|
|||||||
|
|
||||||
fn classify_transition_candidates(
|
fn classify_transition_candidates(
|
||||||
candidates: TransitionCandidateVersions,
|
candidates: TransitionCandidateVersions,
|
||||||
bucket_versioning: RemoteBucketVersioning,
|
initial_versioning: RemoteBucketVersioning,
|
||||||
|
confirmed_versioning: RemoteBucketVersioning,
|
||||||
) -> Result<TransitionCandidateProbe, std::io::Error> {
|
) -> Result<TransitionCandidateProbe, std::io::Error> {
|
||||||
let probe = candidates.classify(bucket_versioning);
|
// GetBucketVersioning and ListObjectVersions are separate requests. An
|
||||||
|
// observed state change makes the combined proof unsafe to persist.
|
||||||
|
if initial_versioning != confirmed_versioning {
|
||||||
|
return Ok(TransitionCandidateProbe::Ambiguous);
|
||||||
|
}
|
||||||
|
let probe = candidates.classify(initial_versioning);
|
||||||
if let TransitionCandidateProbe::VersionedPresent(version_id) = &probe {
|
if let TransitionCandidateProbe::VersionedPresent(version_id) = &probe {
|
||||||
validate_remote_version_id(version_id)?;
|
validate_remote_version_id(version_id)?;
|
||||||
}
|
}
|
||||||
@@ -366,8 +374,19 @@ impl TransitionCandidateVersions {
|
|||||||
};
|
};
|
||||||
|
|
||||||
match bucket_versioning {
|
match bucket_versioning {
|
||||||
RemoteBucketVersioning::Disabled => TransitionCandidateProbe::UnversionedPresent,
|
// A never-versioned S3 object may be listed with either an empty
|
||||||
RemoteBucketVersioning::Suspended if version_id == "null" => TransitionCandidateProbe::VersionedPresent(version_id),
|
// version or the provider's `null` sentinel. Any opaque version is
|
||||||
|
// inconsistent with a disabled bucket and must remain ambiguous.
|
||||||
|
RemoteBucketVersioning::Disabled if version_id.is_empty() || version_id == "null" => {
|
||||||
|
TransitionCandidateProbe::UnversionedPresent
|
||||||
|
}
|
||||||
|
RemoteBucketVersioning::Disabled => TransitionCandidateProbe::Ambiguous,
|
||||||
|
// A `null` version survives a later transition back to Enabled, so
|
||||||
|
// the listed identifier, not only the current bucket status, binds
|
||||||
|
// suspended-null request routing.
|
||||||
|
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled if version_id == "null" => {
|
||||||
|
TransitionCandidateProbe::SuspendedNullPresent
|
||||||
|
}
|
||||||
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled if !version_id.is_empty() => {
|
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled if !version_id.is_empty() => {
|
||||||
TransitionCandidateProbe::VersionedPresent(version_id)
|
TransitionCandidateProbe::VersionedPresent(version_id)
|
||||||
}
|
}
|
||||||
@@ -426,7 +445,8 @@ mod tests {
|
|||||||
for page in pages {
|
for page in pages {
|
||||||
candidates.extend("archive/object", page);
|
candidates.extend("archive/object", page);
|
||||||
}
|
}
|
||||||
candidates.classify(bucket_versioning)
|
classify_transition_candidates(candidates, bucket_versioning, bucket_versioning)
|
||||||
|
.expect("fixture version identifiers are valid")
|
||||||
}
|
}
|
||||||
|
|
||||||
fn candidate_identity() -> TransitionCandidateIdentity {
|
fn candidate_identity() -> TransitionCandidateIdentity {
|
||||||
@@ -505,7 +525,35 @@ mod tests {
|
|||||||
RemoteBucketVersioning::Suspended,
|
RemoteBucketVersioning::Suspended,
|
||||||
&[list_versions(&[("archive/object", "null")], &[], false)],
|
&[list_versions(&[("archive/object", "null")], &[], false)],
|
||||||
),
|
),
|
||||||
TransitionCandidateProbe::VersionedPresent("null".to_string())
|
TransitionCandidateProbe::SuspendedNullPresent
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
classify_pages(
|
||||||
|
RemoteBucketVersioning::Enabled,
|
||||||
|
&[list_versions(&[("archive/object", "null")], &[], false)],
|
||||||
|
),
|
||||||
|
TransitionCandidateProbe::SuspendedNullPresent
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
classify_pages(
|
||||||
|
RemoteBucketVersioning::Suspended,
|
||||||
|
&[list_versions(&[("archive/object", "version-a")], &[], false)],
|
||||||
|
),
|
||||||
|
TransitionCandidateProbe::VersionedPresent("version-a".to_string())
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
classify_pages(
|
||||||
|
RemoteBucketVersioning::Disabled,
|
||||||
|
&[list_versions(&[("archive/object", "null")], &[], false)],
|
||||||
|
),
|
||||||
|
TransitionCandidateProbe::UnversionedPresent
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
classify_pages(
|
||||||
|
RemoteBucketVersioning::Disabled,
|
||||||
|
&[list_versions(&[("archive/object", "unexpected-version")], &[], false)],
|
||||||
|
),
|
||||||
|
TransitionCandidateProbe::Ambiguous
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
classify_pages(RemoteBucketVersioning::Enabled, &[list_versions(&[("archive/object", "")], &[], false)],),
|
classify_pages(RemoteBucketVersioning::Enabled, &[list_versions(&[("archive/object", "")], &[], false)],),
|
||||||
@@ -573,11 +621,23 @@ mod tests {
|
|||||||
let mut candidates = TransitionCandidateVersions::default();
|
let mut candidates = TransitionCandidateVersions::default();
|
||||||
candidates.extend("archive/object", &list_versions(&[("archive/object", "version\ninjection")], &[], false));
|
candidates.extend("archive/object", &list_versions(&[("archive/object", "version\ninjection")], &[], false));
|
||||||
|
|
||||||
let err = classify_transition_candidates(candidates, RemoteBucketVersioning::Enabled)
|
let err = classify_transition_candidates(candidates, RemoteBucketVersioning::Enabled, RemoteBucketVersioning::Enabled)
|
||||||
.expect_err("control characters in listed version IDs must fail closed");
|
.expect_err("control characters in listed version IDs must fail closed");
|
||||||
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn transition_candidate_probe_rejects_bucket_versioning_drift() {
|
||||||
|
let mut candidates = TransitionCandidateVersions::default();
|
||||||
|
candidates.extend("archive/object", &list_versions(&[("archive/object", "null")], &[], false));
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
classify_transition_candidates(candidates, RemoteBucketVersioning::Suspended, RemoteBucketVersioning::Enabled,)
|
||||||
|
.expect("observed versioning drift should be a safe probe result"),
|
||||||
|
TransitionCandidateProbe::Ambiguous
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn remote_bucket_versioning_status_parser_fails_closed() {
|
fn remote_bucket_versioning_status_parser_fails_closed() {
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
|||||||
@@ -3877,6 +3877,31 @@ pub struct SetDisks {
|
|||||||
>,
|
>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Read every available `xl.meta` copy for an exact logical object after
|
||||||
|
/// enforcing this set's metadata read quorum. Callers use the individual
|
||||||
|
/// buffers only for diagnostics and immutable reconciliation digests; a
|
||||||
|
/// missing disk remains visible as `None` and must never be mistaken for an
|
||||||
|
/// absent metadata key.
|
||||||
|
pub(crate) async fn read_legacy_transition_state_metadata_copies(
|
||||||
|
set: &SetDisks,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
) -> Result<Vec<Option<Vec<u8>>>> {
|
||||||
|
let disk_object = rustfs_utils::path::encode_dir_object(object);
|
||||||
|
let disks = set.get_disks_internal().await;
|
||||||
|
if disks.is_empty() {
|
||||||
|
return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object]));
|
||||||
|
}
|
||||||
|
|
||||||
|
let read_quorum = disks.len().div_ceil(2).max(1);
|
||||||
|
let (copies, errs) = SetDisks::read_all_raw_file_info(&disks, bucket, disk_object.as_str(), false).await;
|
||||||
|
if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) {
|
||||||
|
return Err(to_object_err(err.into(), vec![bucket, object]));
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(copies.into_iter().map(|copy| copy.map(|copy| copy.buf)).collect())
|
||||||
|
}
|
||||||
|
|
||||||
// DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones
|
// DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones
|
||||||
// each endpoint's canonical Arc, so an exact Arc set identifies the lock domain.
|
// each endpoint's canonical Arc, so an exact Arc set identifies the lock domain.
|
||||||
pub(crate) fn same_distributed_lock_domain(left: &[Arc<dyn LockClient>], right: &[Arc<dyn LockClient>]) -> bool {
|
pub(crate) fn same_distributed_lock_domain(left: &[Arc<dyn LockClient>], right: &[Arc<dyn LockClient>]) -> bool {
|
||||||
|
|||||||
@@ -11562,6 +11562,7 @@ mod tests {
|
|||||||
pool_index: usize,
|
pool_index: usize,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
object: &str,
|
object: &str,
|
||||||
|
minio_unversioned: bool,
|
||||||
) {
|
) {
|
||||||
for disk_index in 0..4 {
|
for disk_index in 0..4 {
|
||||||
let metadata_path =
|
let metadata_path =
|
||||||
@@ -11595,6 +11596,11 @@ mod tests {
|
|||||||
] {
|
] {
|
||||||
rustfs_utils::http::metadata_compat::remove_bytes(&mut object_meta.meta_sys, suffix);
|
rustfs_utils::http::metadata_compat::remove_bytes(&mut object_meta.meta_sys, suffix);
|
||||||
}
|
}
|
||||||
|
if minio_unversioned {
|
||||||
|
object_meta
|
||||||
|
.meta_sys
|
||||||
|
.insert("x-minio-internal-transitioned-versionID".to_string(), Vec::new());
|
||||||
|
}
|
||||||
*shallow = rustfs_filemeta::FileMetaShallowVersion::try_from(version)
|
*shallow = rustfs_filemeta::FileMetaShallowVersion::try_from(version)
|
||||||
.expect("legacy transitioned version should re-encode");
|
.expect("legacy transitioned version should re-encode");
|
||||||
}
|
}
|
||||||
@@ -11605,6 +11611,152 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
async fn read_store_body(
|
||||||
|
store: &Arc<crate::store::ECStore>,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
range: Option<HTTPRangeSpec>,
|
||||||
|
opts: &ObjectOptions,
|
||||||
|
) -> Vec<u8> {
|
||||||
|
let mut reader = store
|
||||||
|
.get_object_reader(bucket, object, range, HeaderMap::new(), opts)
|
||||||
|
.await
|
||||||
|
.expect("object reader should open");
|
||||||
|
let mut body = Vec::new();
|
||||||
|
reader.stream.read_to_end(&mut body).await.expect("object body should drain");
|
||||||
|
body
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
async fn legacy_unknown_unversioned_transition_supports_head_get_and_range_without_backfill() {
|
||||||
|
let temp_dir = tempfile::tempdir().expect("create legacy unknown unversioned store dir");
|
||||||
|
let (ctx, store, _shutdown) =
|
||||||
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "legacy-unknown-unversioned-read", &[4])).await;
|
||||||
|
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||||
|
let tier_name = "LEGACY-UNKNOWN-UNVERSIONED-READ";
|
||||||
|
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
|
||||||
|
backend.set_put_remote_version(Some(String::new())).await;
|
||||||
|
let bucket = "legacy-unknown-unversioned-read-bucket";
|
||||||
|
let object = "object.bin";
|
||||||
|
let payload = b"legacy unversioned remote tier object remains readable".repeat(1024);
|
||||||
|
store
|
||||||
|
.make_bucket(bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("legacy source bucket should be created");
|
||||||
|
let mut reader = PutObjReader::from_vec(payload.clone());
|
||||||
|
let source = store
|
||||||
|
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("legacy source should be written");
|
||||||
|
store
|
||||||
|
.transition_object(
|
||||||
|
bucket,
|
||||||
|
object,
|
||||||
|
&ObjectOptions {
|
||||||
|
transition: TransitionOptions {
|
||||||
|
status: TRANSITION_PENDING.to_string(),
|
||||||
|
tier: tier_name.to_string(),
|
||||||
|
etag: source.etag.clone().expect("legacy source should have an etag"),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
mod_time: source.mod_time,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("legacy source should transition");
|
||||||
|
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, true).await;
|
||||||
|
backend.clear_op_log().await;
|
||||||
|
|
||||||
|
let opts = ObjectOptions {
|
||||||
|
metadata_cache_safe: false,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let head = store
|
||||||
|
.get_object_info(bucket, object, &opts)
|
||||||
|
.await
|
||||||
|
.expect("legacy transitioned HEAD should use local metadata");
|
||||||
|
assert_eq!(head.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown);
|
||||||
|
assert!(head.transitioned_object.version_id.is_empty());
|
||||||
|
assert_eq!(
|
||||||
|
head.user_defined
|
||||||
|
.get("x-minio-internal-transitioned-versionID")
|
||||||
|
.map(String::as_str),
|
||||||
|
Some(""),
|
||||||
|
"the MinIO empty version-key provenance must survive xl.meta decoding"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!rustfs_utils::http::metadata_compat::contains_key_str(
|
||||||
|
&head.user_defined,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
),
|
||||||
|
"the compatibility read must not synthesize version-state metadata"
|
||||||
|
);
|
||||||
|
|
||||||
|
let full_body = read_store_body(&store, bucket, object, None, &opts).await;
|
||||||
|
assert_eq!(full_body, payload);
|
||||||
|
|
||||||
|
let range = HTTPRangeSpec {
|
||||||
|
is_suffix_length: false,
|
||||||
|
start: 7,
|
||||||
|
end: 38,
|
||||||
|
};
|
||||||
|
let ranged_body = read_store_body(&store, bucket, object, Some(range), &opts).await;
|
||||||
|
assert_eq!(ranged_body, &payload[7..=38]);
|
||||||
|
|
||||||
|
let after_read = store.pools[0]
|
||||||
|
.get_disks_by_key(object)
|
||||||
|
.load_file_info_versions_exact(bucket, object)
|
||||||
|
.await
|
||||||
|
.expect("legacy metadata should remain readable after GET")
|
||||||
|
.expect("legacy object metadata should remain on disk")
|
||||||
|
.versions
|
||||||
|
.into_iter()
|
||||||
|
.find(|version| version.transition_status == rustfs_filemeta::TRANSITION_COMPLETE)
|
||||||
|
.expect("legacy transitioned source should remain visible after GET");
|
||||||
|
assert_eq!(after_read.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown);
|
||||||
|
assert!(after_read.transition_version.is_none());
|
||||||
|
assert!(after_read.transition_version_id.is_none());
|
||||||
|
assert_eq!(
|
||||||
|
after_read
|
||||||
|
.metadata
|
||||||
|
.get("x-minio-internal-transitioned-versionID")
|
||||||
|
.map(String::as_str),
|
||||||
|
Some(""),
|
||||||
|
"the MinIO empty version-key provenance must remain after GET and Range GET"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!rustfs_utils::http::metadata_compat::contains_key_str(
|
||||||
|
&after_read.metadata,
|
||||||
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
),
|
||||||
|
"the compatibility read must remain side-effect free"
|
||||||
|
);
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
backend.op_log().await,
|
||||||
|
vec![
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: after_read.transitioned_objname.clone(),
|
||||||
|
},
|
||||||
|
MockWarmOp::Get {
|
||||||
|
object: after_read.transitioned_objname.clone(),
|
||||||
|
},
|
||||||
|
MockWarmOp::Probe {
|
||||||
|
object: after_read.transitioned_objname.clone(),
|
||||||
|
},
|
||||||
|
MockWarmOp::Get {
|
||||||
|
object: after_read.transitioned_objname,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
"legacy reads should probe before each unversioned GET and never mutate local metadata"
|
||||||
|
);
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial(storage_class_env)]
|
#[serial_test::serial(storage_class_env)]
|
||||||
@@ -11645,7 +11797,7 @@ mod tests {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("legacy source should transition");
|
.expect("legacy source should transition");
|
||||||
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object).await;
|
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, false).await;
|
||||||
let legacy = store.pools[0]
|
let legacy = store.pools[0]
|
||||||
.get_disks_by_key(object)
|
.get_disks_by_key(object)
|
||||||
.load_file_info_versions_exact(bucket, object)
|
.load_file_info_versions_exact(bucket, object)
|
||||||
@@ -12786,7 +12938,7 @@ mod tests {
|
|||||||
.expect("merge-loser source should transition");
|
.expect("merge-loser source should transition");
|
||||||
copy_test_xlmeta_between_pools(temp_dir.path(), 0, 1, bucket, object).await;
|
copy_test_xlmeta_between_pools(temp_dir.path(), 0, 1, bucket, object).await;
|
||||||
}
|
}
|
||||||
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin").await;
|
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin", false).await;
|
||||||
backend.set_remove_failure(true);
|
backend.set_remove_failure(true);
|
||||||
store.pools[1]
|
store.pools[1]
|
||||||
.delete_object(bucket, "hidden/item.bin", ObjectOptions::default())
|
.delete_object(bucket, "hidden/item.bin", ObjectOptions::default())
|
||||||
|
|||||||
@@ -297,6 +297,20 @@ fn transitioned_version_from_bytes(value: Option<&[u8]>, state: TransitionVersio
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn transition_version_metadata_value(raw: &[u8], decoded: Option<&str>) -> String {
|
||||||
|
decoded.map(str::to_owned).unwrap_or_else(|| {
|
||||||
|
if raw.is_empty() {
|
||||||
|
String::new()
|
||||||
|
} else {
|
||||||
|
String::from_utf8_lossy(raw).into_owned()
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
fn is_transition_version_metadata_key(key: &str) -> bool {
|
||||||
|
strip_internal_prefix_preserving_case(key).is_some_and(|suffix| suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID))
|
||||||
|
}
|
||||||
|
|
||||||
fn validate_transition_version_state(state: TransitionVersionState, version: Option<&str>) -> Result<()> {
|
fn validate_transition_version_state(state: TransitionVersionState, version: Option<&str>) -> Result<()> {
|
||||||
let valid = match state {
|
let valid = match state {
|
||||||
TransitionVersionState::Unknown | TransitionVersionState::KnownDisabled => version.is_none(),
|
TransitionVersionState::Unknown | TransitionVersionState::KnownDisabled => version.is_none(),
|
||||||
@@ -366,14 +380,26 @@ impl<'a> DerivedInternalMetadata<'a> {
|
|||||||
}
|
}
|
||||||
*slot = Some(value.as_slice());
|
*slot = Some(value.as_slice());
|
||||||
}
|
}
|
||||||
|
fn merge_consistent<'a>(canonical: Option<&'a [u8]>, legacy: Option<&'a [u8]>) -> Result<Option<&'a [u8]>> {
|
||||||
|
if let (Some(canonical), Some(legacy)) = (canonical, legacy)
|
||||||
|
&& canonical != legacy
|
||||||
|
{
|
||||||
|
return Err(Error::FileCorrupt);
|
||||||
|
}
|
||||||
|
Ok(canonical.or(legacy))
|
||||||
|
}
|
||||||
|
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
checksum: canonical.checksum.or(legacy.checksum),
|
checksum: canonical.checksum.or(legacy.checksum),
|
||||||
part_checksums: canonical.part_checksums.or(legacy.part_checksums),
|
part_checksums: canonical.part_checksums.or(legacy.part_checksums),
|
||||||
transition_status: canonical.transition_status.or(legacy.transition_status),
|
transition_status: merge_consistent(canonical.transition_status, legacy.transition_status)?,
|
||||||
transitioned_object: canonical.transitioned_object.or(legacy.transitioned_object),
|
transitioned_object: merge_consistent(canonical.transitioned_object, legacy.transitioned_object)?,
|
||||||
transitioned_version: canonical.transitioned_version.or(legacy.transitioned_version),
|
transitioned_version: merge_consistent(canonical.transitioned_version, legacy.transitioned_version)?,
|
||||||
transitioned_version_state: canonical.transitioned_version_state.or(legacy.transitioned_version_state),
|
transitioned_version_state: merge_consistent(
|
||||||
transition_tier: canonical.transition_tier.or(legacy.transition_tier),
|
canonical.transitioned_version_state,
|
||||||
|
legacy.transitioned_version_state,
|
||||||
|
)?,
|
||||||
|
transition_tier: merge_consistent(canonical.transition_tier, legacy.transition_tier)?,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -438,8 +464,14 @@ impl FileInfo {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn set_transition_version_state(meta_sys: &mut HashMap<String, Vec<u8>>, state: TransitionVersionState) {
|
fn set_transition_version_state(
|
||||||
if state == TransitionVersionState::Unknown {
|
meta_sys: &mut HashMap<String, Vec<u8>>,
|
||||||
|
state: TransitionVersionState,
|
||||||
|
source_metadata: &HashMap<String, String>,
|
||||||
|
) {
|
||||||
|
if state == TransitionVersionState::Unknown
|
||||||
|
&& !rustfs_utils::http::metadata_compat::contains_key_str(source_metadata, SUFFIX_TRANSITIONED_VERSION_STATE)
|
||||||
|
{
|
||||||
remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE);
|
remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE);
|
||||||
} else {
|
} else {
|
||||||
insert_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str().as_bytes().to_vec());
|
insert_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str().as_bytes().to_vec());
|
||||||
@@ -2643,6 +2675,11 @@ impl MetaObject {
|
|||||||
if derived_metadata.transitioned_version_state.is_some() {
|
if derived_metadata.transitioned_version_state.is_some() {
|
||||||
validate_transition_version_state(transition_version_state, transition_version.as_deref())?;
|
validate_transition_version_state(transition_version_state, transition_version.as_deref())?;
|
||||||
}
|
}
|
||||||
|
for (key, value) in &self.meta_sys {
|
||||||
|
if is_transition_version_metadata_key(key) {
|
||||||
|
metadata.insert(key.to_owned(), transition_version_metadata_value(value, transition_version.as_deref()));
|
||||||
|
}
|
||||||
|
}
|
||||||
let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok());
|
let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok());
|
||||||
let transition_tier = derived_metadata
|
let transition_tier = derived_metadata
|
||||||
.transition_tier
|
.transition_tier
|
||||||
@@ -2689,7 +2726,7 @@ impl MetaObject {
|
|||||||
} else {
|
} else {
|
||||||
remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID);
|
remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID);
|
||||||
}
|
}
|
||||||
set_transition_version_state(&mut self.meta_sys, fi.transition_version_state);
|
set_transition_version_state(&mut self.meta_sys, fi.transition_version_state, &fi.metadata);
|
||||||
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec());
|
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec());
|
||||||
if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) {
|
if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) {
|
||||||
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes());
|
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes());
|
||||||
@@ -2830,7 +2867,7 @@ impl From<FileInfo> for MetaObject {
|
|||||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
||||||
}
|
}
|
||||||
if !value.transition_status.is_empty() {
|
if !value.transition_status.is_empty() {
|
||||||
set_transition_version_state(&mut meta_sys, value.transition_version_state);
|
set_transition_version_state(&mut meta_sys, value.transition_version_state, &value.metadata);
|
||||||
}
|
}
|
||||||
|
|
||||||
if !value.transition_tier.is_empty() {
|
if !value.transition_tier.is_empty() {
|
||||||
@@ -2985,6 +3022,12 @@ impl MetaDeleteMarker {
|
|||||||
fi.transition_version_state = transition_version_state_from_bytes(derived_metadata.transitioned_version_state)?;
|
fi.transition_version_state = transition_version_state_from_bytes(derived_metadata.transitioned_version_state)?;
|
||||||
fi.transition_version =
|
fi.transition_version =
|
||||||
transitioned_version_from_bytes(derived_metadata.transitioned_version, fi.transition_version_state);
|
transitioned_version_from_bytes(derived_metadata.transitioned_version, fi.transition_version_state);
|
||||||
|
for (key, value) in &self.meta_sys {
|
||||||
|
if is_transition_version_metadata_key(key) {
|
||||||
|
fi.metadata
|
||||||
|
.insert(key.to_owned(), transition_version_metadata_value(value, fi.transition_version.as_deref()));
|
||||||
|
}
|
||||||
|
}
|
||||||
fi.transition_version_id = fi.transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok());
|
fi.transition_version_id = fi.transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok());
|
||||||
if derived_metadata.transitioned_version_state.is_some() {
|
if derived_metadata.transitioned_version_state.is_some() {
|
||||||
validate_transition_version_state(fi.transition_version_state, fi.transition_version.as_deref())?;
|
validate_transition_version_state(fi.transition_version_state, fi.transition_version.as_deref())?;
|
||||||
@@ -3152,7 +3195,7 @@ impl From<FileInfo> for MetaDeleteMarker {
|
|||||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
||||||
}
|
}
|
||||||
if !value.transition_status.is_empty() || value.tier_free_version() {
|
if !value.transition_status.is_empty() || value.tier_free_version() {
|
||||||
set_transition_version_state(&mut meta_sys, value.transition_version_state);
|
set_transition_version_state(&mut meta_sys, value.transition_version_state, &value.metadata);
|
||||||
}
|
}
|
||||||
if !value.transition_tier.is_empty() {
|
if !value.transition_tier.is_empty() {
|
||||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec());
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec());
|
||||||
@@ -4574,6 +4617,7 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false)
|
.into_fileinfo("b", "k", false)
|
||||||
.expect("into_fileinfo");
|
.expect("into_fileinfo");
|
||||||
assert_eq!(fi.transition_version_id, None);
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(String::new()));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4585,6 +4629,10 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false)
|
.into_fileinfo("b", "k", false)
|
||||||
.expect("into_fileinfo");
|
.expect("into_fileinfo");
|
||||||
assert_eq!(fi.transition_version_id, None);
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert!(
|
||||||
|
get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()),
|
||||||
|
"nil UUID bytes must remain distinguishable from an empty MinIO version"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4598,6 +4646,7 @@ mod tests {
|
|||||||
assert_eq!(fi.transition_version_id, Some(id));
|
assert_eq!(fi.transition_version_id, Some(id));
|
||||||
assert_eq!(fi.transition_version, Some(id.to_string()));
|
assert_eq!(fi.transition_version, Some(id.to_string()));
|
||||||
assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown);
|
assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown);
|
||||||
|
assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(id.to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4637,6 +4686,36 @@ mod tests {
|
|||||||
assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown);
|
assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn meta_object_transition_version_state_explicit_unknown_is_not_legacy_missing() {
|
||||||
|
let mut metadata = HashMap::new();
|
||||||
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
|
&mut metadata,
|
||||||
|
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
TransitionVersionState::Unknown.as_str().to_string(),
|
||||||
|
);
|
||||||
|
let fi = FileInfo {
|
||||||
|
transition_status: "complete".to_string(),
|
||||||
|
transition_version_state: TransitionVersionState::Unknown,
|
||||||
|
metadata,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let object = MetaObject::from(fi);
|
||||||
|
assert_eq!(
|
||||||
|
get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE),
|
||||||
|
Some(b"unknown".as_slice())
|
||||||
|
);
|
||||||
|
let decoded = object
|
||||||
|
.into_fileinfo("b", "k", false)
|
||||||
|
.expect("explicit unknown state should decode");
|
||||||
|
assert_eq!(decoded.transition_version_state, TransitionVersionState::Unknown);
|
||||||
|
assert_eq!(
|
||||||
|
rustfs_utils::http::metadata_compat::get_consistent_str(&decoded.metadata, SUFFIX_TRANSITIONED_VERSION_STATE,),
|
||||||
|
Some("unknown")
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn meta_object_transition_version_state_exact_round_trips_dual_keys() {
|
fn meta_object_transition_version_state_exact_round_trips_dual_keys() {
|
||||||
let id = sample_version_id();
|
let id = sample_version_id();
|
||||||
@@ -4753,6 +4832,10 @@ mod tests {
|
|||||||
.expect("invalid transition version bytes must not fail the object read");
|
.expect("invalid transition version bytes must not fail the object read");
|
||||||
assert_eq!(fi.transition_version_id, None);
|
assert_eq!(fi.transition_version_id, None);
|
||||||
assert_eq!(fi.transition_version, None);
|
assert_eq!(fi.transition_version, None);
|
||||||
|
assert!(
|
||||||
|
get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()),
|
||||||
|
"invalid raw bytes must remain distinguishable from an empty MinIO version"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4795,6 +4878,10 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false)
|
.into_fileinfo("b", "k", false)
|
||||||
.expect("nil tier version should remain an absent remote version");
|
.expect("nil tier version should remain an absent remote version");
|
||||||
assert_eq!(fi.transition_version_id, None);
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert!(
|
||||||
|
get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()),
|
||||||
|
"nil UUID bytes must remain distinguishable from an empty MinIO version"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4812,6 +4899,7 @@ mod tests {
|
|||||||
.expect("legacy binary UUID tier version should decode");
|
.expect("legacy binary UUID tier version should decode");
|
||||||
assert_eq!(fi.transition_version_id, Some(id));
|
assert_eq!(fi.transition_version_id, Some(id));
|
||||||
assert_eq!(fi.transition_version, Some(id.to_string()));
|
assert_eq!(fi.transition_version, Some(id.to_string()));
|
||||||
|
assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(id.to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4910,6 +4998,23 @@ mod tests {
|
|||||||
assert_eq!(err, Error::FileCorrupt);
|
assert_eq!(err, Error::FileCorrupt);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn meta_object_transition_version_state_mixed_case_alias_conflict_fails_closed() {
|
||||||
|
let sys = HashMap::from([
|
||||||
|
(
|
||||||
|
format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}"),
|
||||||
|
b"unknown".to_vec(),
|
||||||
|
),
|
||||||
|
("X-Minio-Internal-transitioned-version-state".to_string(), b"exact".to_vec()),
|
||||||
|
]);
|
||||||
|
|
||||||
|
let err = make_meta_object_with_sys(sys)
|
||||||
|
.into_fileinfo("b", "k", false)
|
||||||
|
.expect_err("mixed-case transition state aliases must agree");
|
||||||
|
|
||||||
|
assert_eq!(err, Error::FileCorrupt);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn version_header_sorts_before_prefers_object_over_delete_marker_on_equal_mod_time() {
|
fn version_header_sorts_before_prefers_object_over_delete_marker_on_equal_mod_time() {
|
||||||
let object = FileMetaVersionHeader {
|
let object = FileMetaVersionHeader {
|
||||||
|
|||||||
@@ -15,12 +15,13 @@
|
|||||||
use crate::admin::auth::authorize_admin_request;
|
use crate::admin::auth::authorize_admin_request;
|
||||||
use crate::admin::router::{AdminOperation, Operation, S3Router};
|
use crate::admin::router::{AdminOperation, Operation, S3Router};
|
||||||
use crate::admin::runtime_sources::object_store_from_extensions;
|
use crate::admin::runtime_sources::object_store_from_extensions;
|
||||||
use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket;
|
use crate::admin::storage_api::bucket::{is_reserved_or_invalid_bucket, utils::is_valid_object_prefix};
|
||||||
use crate::admin::storage_api::error::StorageError;
|
use crate::admin::storage_api::error::StorageError;
|
||||||
use crate::admin::storage_api::lifecycle::{
|
use crate::admin::storage_api::lifecycle::{
|
||||||
ManualTransitionCancelCheck, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionProgressSink,
|
LegacyTransitionStateReconcileError, LegacyTransitionStateReconcileRequest, LegacyTransitionStateReconcileResponse,
|
||||||
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, ManualTransitionScopeAdmission,
|
LegacyTransitionStateReconcileSelector, ManualTransitionCancelCheck, ManualTransitionJobRecord, ManualTransitionJobState,
|
||||||
ManualTransitionScopeAdmissionClaim, TransitionOperatorDeleteResult, TransitionOperatorError,
|
ManualTransitionProgressSink, ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport,
|
||||||
|
ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim, TransitionOperatorDeleteResult, TransitionOperatorError,
|
||||||
claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
|
claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
|
||||||
delete_transition_candidate_for_operator, enqueue_transition_for_existing_objects_scoped,
|
delete_transition_candidate_for_operator, enqueue_transition_for_existing_objects_scoped,
|
||||||
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
|
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
|
||||||
@@ -230,6 +231,16 @@ pub fn register_ilm_transition_route(r: &mut S3Router<AdminOperation>) -> std::i
|
|||||||
format!("{ADMIN_PREFIX}/v3/ilm/transition/reconcile/{{transaction_id}}").as_str(),
|
format!("{ADMIN_PREFIX}/v3/ilm/transition/reconcile/{{transaction_id}}").as_str(),
|
||||||
AdminOperation(&TransitionReconcileApplyHandler {}),
|
AdminOperation(&TransitionReconcileApplyHandler {}),
|
||||||
)?;
|
)?;
|
||||||
|
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -423,6 +434,92 @@ fn transition_transaction_id_from_params(params: &Params<'_, '_>) -> S3Result<Uu
|
|||||||
.map_err(|_| s3_error!(InvalidArgument, "invalid transition transaction id"))
|
.map_err(|_| s3_error!(InvalidArgument, "invalid transition transaction id"))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Deserialize)]
|
||||||
|
#[serde(deny_unknown_fields)]
|
||||||
|
struct LegacyTransitionStateReconcileQuery {
|
||||||
|
bucket: Option<String>,
|
||||||
|
object: Option<String>,
|
||||||
|
#[serde(rename = "versionId")]
|
||||||
|
version_id: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn parse_legacy_transition_state_reconcile_query(query: Option<&str>) -> S3Result<LegacyTransitionStateReconcileSelector> {
|
||||||
|
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 {
|
fn map_transition_operator_error(err: TransitionOperatorError) -> S3Error {
|
||||||
match err {
|
match err {
|
||||||
TransitionOperatorError::NotFound => s3_error!(NoSuchKey, "transition transaction not found"),
|
TransitionOperatorError::NotFound => s3_error!(NoSuchKey, "transition transaction not found"),
|
||||||
@@ -1089,6 +1186,52 @@ impl Operation for TransitionReconcileApplyHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub struct LegacyTransitionStateReconcileInspectHandler {}
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl Operation for LegacyTransitionStateReconcileInspectHandler {
|
||||||
|
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||||
|
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<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||||
|
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)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
@@ -1171,12 +1314,87 @@ mod tests {
|
|||||||
let apply = src
|
let apply = src
|
||||||
.split("impl Operation for TransitionReconcileApplyHandler")
|
.split("impl Operation for TransitionReconcileApplyHandler")
|
||||||
.nth(1)
|
.nth(1)
|
||||||
.and_then(|block| block.split("#[cfg(test)]").next())
|
.and_then(|block| block.split("pub struct LegacyTransitionStateReconcileInspectHandler").next())
|
||||||
.expect("apply handler block");
|
.expect("apply handler block");
|
||||||
assert!(apply.contains("AdminAction::SetTierAction"));
|
assert!(apply.contains("AdminAction::SetTierAction"));
|
||||||
assert!(!apply.contains("AdminAction::ListTierAction"));
|
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]
|
#[test]
|
||||||
fn manual_transition_query_defaults_to_bounded_run() {
|
fn manual_transition_query_defaults_to_bounded_run() {
|
||||||
let (bucket, options, run_mode) =
|
let (bucket, options, run_mode) =
|
||||||
@@ -1632,6 +1850,24 @@ mod tests {
|
|||||||
assert_eq!(cancel_err.message(), Some("authentication required"));
|
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(manual_transition_job_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(manual_transition_job_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]
|
#[test]
|
||||||
fn manual_transition_job_handlers_authorize_validate_and_load_store() {
|
fn manual_transition_job_handlers_authorize_validate_and_load_store() {
|
||||||
let src = include_str!("ilm_transition.rs");
|
let src = include_str!("ilm_transition.rs");
|
||||||
|
|||||||
@@ -520,6 +520,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
|
|||||||
SET_TIER,
|
SET_TIER,
|
||||||
RouteRiskLevel::High,
|
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(
|
admin(
|
||||||
HttpMethod::Get,
|
HttpMethod::Get,
|
||||||
"/rustfs/admin/v3/audit/target/list",
|
"/rustfs/admin/v3/audit/target/list",
|
||||||
@@ -2164,11 +2176,15 @@ mod tests {
|
|||||||
assert_action(HttpMethod::Delete, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", SET_TIER);
|
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::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::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::Post, "/rustfs/admin/v3/ilm/transition/run", SERVER_INFO);
|
||||||
assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/transition/jobs/{job_id}", 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::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::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::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]
|
#[test]
|
||||||
|
|||||||
@@ -225,6 +225,8 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
|
|||||||
"/v3/ilm/transition/reconcile/{transaction_id}",
|
"/v3/ilm/transition/reconcile/{transaction_id}",
|
||||||
"/v3/ilm/transition/reconcile/11111111-1111-4111-8111-111111111111",
|
"/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(
|
admin_route_sample(
|
||||||
Method::DELETE,
|
Method::DELETE,
|
||||||
"/v3/ilm/transition/jobs/{job_id}",
|
"/v3/ilm/transition/jobs/{job_id}",
|
||||||
@@ -951,6 +953,8 @@ fn test_register_routes_cover_representative_admin_paths() {
|
|||||||
Method::POST,
|
Method::POST,
|
||||||
&admin_path("/v3/ilm/transition/reconcile/11111111-1111-4111-8111-111111111111"),
|
&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::GET, &table_catalog_path("/config"));
|
||||||
assert_route(&router, Method::PUT, &table_catalog_path("/buckets/analytics"));
|
assert_route(&router, Method::PUT, &table_catalog_path("/buckets/analytics"));
|
||||||
|
|||||||
@@ -223,6 +223,10 @@ pub(crate) mod lifecycle {
|
|||||||
renew_manual_transition_job_lease_if_owned, request_manual_transition_job_cancel, save_manual_transition_job_record,
|
renew_manual_transition_job_lease_if_owned, request_manual_transition_job_cancel, save_manual_transition_job_record,
|
||||||
update_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 =
|
pub(crate) type ManualTransitionCancelCheck =
|
||||||
super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionCancelCheck;
|
super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionCancelCheck;
|
||||||
pub(crate) type ManualTransitionQueueSnapshot =
|
pub(crate) type ManualTransitionQueueSnapshot =
|
||||||
|
|||||||
Reference in New Issue
Block a user