mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-21 03:46:37 +00:00
fix(tiering): replay exact cleanup journals
This commit is contained in:
@@ -10194,6 +10194,61 @@ mod tests {
|
|||||||
assert_eq!(backend.remove_count().await, 0);
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn journal_replay_deletes_confirmed_exact_provider_token() {
|
||||||
|
let (_disk_paths, ecstore) = setup_test_env().await;
|
||||||
|
let (backend, _) = register_recovery_mock_tier(&ecstore).await;
|
||||||
|
let lease = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM")
|
||||||
|
.await
|
||||||
|
.expect("mock tier lease should be available");
|
||||||
|
let identity = lease.backend_identity();
|
||||||
|
backend
|
||||||
|
.set_put_remote_version(Some("provider-version-token".to_string()))
|
||||||
|
.await;
|
||||||
|
lease
|
||||||
|
.put(
|
||||||
|
"remote/object",
|
||||||
|
crate::client::transition_api::ReaderImpl::Body(bytes::Bytes::from_static(b"candidate")),
|
||||||
|
9,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("confirmed remote candidate should be seeded");
|
||||||
|
backend.set_remove_failure(true);
|
||||||
|
backend.set_reject_non_empty_remote_versions(true);
|
||||||
|
let je = Jentry {
|
||||||
|
obj_name: "remote/object".to_string(),
|
||||||
|
version_id: "provider-version-token".to_string(),
|
||||||
|
tier_name: "WARM".to_string(),
|
||||||
|
backend_identity: Some(identity),
|
||||||
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
|
};
|
||||||
|
|
||||||
|
crate::set_disk::cleanup_rejected_transition_upload_durably(
|
||||||
|
&lease,
|
||||||
|
&je.obj_name,
|
||||||
|
&je.version_id,
|
||||||
|
true,
|
||||||
|
Some(ecstore.clone()),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("failed immediate cleanup should remain durable in the journal");
|
||||||
|
assert!(backend.contains(&je.obj_name).await);
|
||||||
|
|
||||||
|
backend.set_remove_failure(false);
|
||||||
|
crate::bucket::lifecycle::tier_delete_journal::process_tier_delete_journal_entry(ecstore, &je)
|
||||||
|
.await
|
||||||
|
.expect("identity-bound exact journal must retry confirmed candidate cleanup");
|
||||||
|
|
||||||
|
assert!(!backend.contains(&je.obj_name).await);
|
||||||
|
assert_eq!(backend.exact_remove_count(), 2);
|
||||||
|
assert_eq!(
|
||||||
|
backend.remove_versions().await,
|
||||||
|
vec![("remote/object".to_string(), "provider-version-token".to_string())]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
async fn seed_recoverable_free_version(
|
async fn seed_recoverable_free_version(
|
||||||
disk_paths: &[PathBuf],
|
disk_paths: &[PathBuf],
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
|
|||||||
@@ -20,7 +20,10 @@ use tokio_util::sync::CancellationToken;
|
|||||||
use tracing::{debug, warn};
|
use tracing::{debug, warn};
|
||||||
|
|
||||||
use crate::bucket::lifecycle::config_boundary;
|
use crate::bucket::lifecycle::config_boundary;
|
||||||
use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent_with_manager_and_identity};
|
use crate::bucket::lifecycle::tier_sweeper::{
|
||||||
|
Jentry, delete_confirmed_transition_candidate_exact_with_manager_and_identity,
|
||||||
|
delete_object_from_remote_tier_idempotent_with_manager_and_identity,
|
||||||
|
};
|
||||||
use crate::disk::RUSTFS_META_BUCKET;
|
use crate::disk::RUSTFS_META_BUCKET;
|
||||||
use crate::error::{Error, Result};
|
use crate::error::{Error, Result};
|
||||||
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
|
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
|
||||||
@@ -270,15 +273,26 @@ pub async fn process_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jentry) -
|
|||||||
let backend_identity = je
|
let backend_identity = je
|
||||||
.backend_identity
|
.backend_identity
|
||||||
.ok_or_else(|| std::io::Error::other("legacy tier delete journal has no durable backend identity"))?;
|
.ok_or_else(|| std::io::Error::other("legacy tier delete journal has no durable backend identity"))?;
|
||||||
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
if je.version_id_exact {
|
||||||
&je.obj_name,
|
delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
&je.version_id,
|
&je.obj_name,
|
||||||
&je.tier_name,
|
&je.version_id,
|
||||||
backend_identity,
|
&je.tier_name,
|
||||||
&api.tier_config_mgr(),
|
backend_identity,
|
||||||
je.version_id_exact,
|
&api.tier_config_mgr(),
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
} else {
|
||||||
|
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
||||||
|
&je.obj_name,
|
||||||
|
&je.version_id,
|
||||||
|
&je.tier_name,
|
||||||
|
backend_identity,
|
||||||
|
&api.tier_config_mgr(),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
remove_tier_delete_journal_entry(api, je).await
|
remove_tier_delete_journal_entry(api, je).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -793,6 +793,18 @@ mod test {
|
|||||||
backend.remove_versions().await,
|
backend.remove_versions().await,
|
||||||
vec![("remote/object".to_string(), "provider-version-token".to_string())]
|
vec![("remote/object".to_string(), "provider-version-token".to_string())]
|
||||||
);
|
);
|
||||||
|
|
||||||
|
let err = delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
|
"remote/object",
|
||||||
|
"",
|
||||||
|
"WARM",
|
||||||
|
identity,
|
||||||
|
&manager,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect_err("confirmed versioned cleanup must reject an empty token");
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
|
||||||
|
assert_eq!(backend.remove_count().await, 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
@@ -690,6 +690,8 @@ mod ops;
|
|||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
||||||
pub(crate) use ops::object::body_cache_plaintext_len;
|
pub(crate) use ops::object::body_cache_plaintext_len;
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) use ops::object::cleanup_rejected_transition_upload_durably;
|
||||||
mod read;
|
mod read;
|
||||||
mod replication;
|
mod replication;
|
||||||
pub(crate) mod shard_source;
|
pub(crate) mod shard_source;
|
||||||
|
|||||||
@@ -1807,7 +1807,7 @@ impl Drop for TransitionUploadCleanup {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn cleanup_rejected_transition_upload_durably(
|
pub(crate) async fn cleanup_rejected_transition_upload_durably(
|
||||||
lease: &TierOperationLease,
|
lease: &TierOperationLease,
|
||||||
object: &str,
|
object: &str,
|
||||||
cleanup_version: &str,
|
cleanup_version: &str,
|
||||||
|
|||||||
Reference in New Issue
Block a user