fix(lifecycle): safely expire all object versions (#6291)

* fix(lifecycle): safely expire all object versions

* fix(lifecycle): preserve delete-all replication purges

* fix(lifecycle): remove dead replication journal

* fix(ci): avoid lifecycle transition test stack overflow

* fix(lifecycle): release recovery locks before tier IO

* test(lifecycle): align object-lock error assertions

* test(lifecycle): avoid scanner restore stack overflow

* test(scanner): avoid stack overflow in transition and restore flow test (#6300)

* refactor(scanner): split remote_scanner.rs into stream child module (#6289)

Split the 3080-line remote_scanner.rs (47% inline tests) into a
canonical foo.rs + foo/ module tree with zero behavior change:

- remote_scanner.rs (~320): protocol constants, process statics, and
  the request decode/validate/admit/preflight/claim API plus root
  re-exports
- remote_scanner/stream.rs (~1340): wire/frame types, replay cache,
  FrameAuthenticator, serve path, local bucket scan + persist, client
  scan, and the bounded stream plumbing
- remote_scanner/stream/tests.rs (~1470): the inline test module as a
  child module of stream so it can reach both parents' private items

All crate paths are unchanged: lib.rs re-exports
(serve_remote_scanner_request, RemoteScannerRequest, ...) resolve
through root re-exports, and scanner_io's crate::remote_scanner::
{scan_remote_bucket, RemoteScannerScanSpec, RemoteScannerOutcome}
paths resolve through pub(crate) re-exports. Cross-module items gain
pub(super), whose scope equals the old single-module privacy domain;
no item's effective visibility widens. Code is moved verbatim apart
from those markers, per-module import headers, and rustfmt line
re-wraps.

Co-authored-by: heihutu <heihutu@gmail.com>

* refactor(heal): split resume.rs into focused child modules (#6290)

Split the 4242-line resume.rs (46% inline tests) into a canonical
foo.rs + foo/ module tree with zero behavior change:

- resume.rs (~1020): state file constants, PersistThrottle, ResumeState,
  ResumeManager core (constructors, load/discovery, progress mutators,
  ordinary persistence) plus root re-exports
- resume/replacement.rs (~690): replacement-intent/proof types and the
  ResumeManager replacement-lifecycle methods
- resume/checkpoint.rs (~350): ResumeCheckpoint + CheckpointManager
- resume/utils.rs (~310): ResumeUtils statics
- resume/tests.rs (~1980): the inline test module as a child module

All module paths are unchanged (heal::resume::CheckpointManager and
friends resolve through root re-exports), so no consumer inside or
outside the crate changes. Items defined in child modules keep
module-private visibility; only the ten cross-module helpers gain
pub(super), which is not part of the crate API. Code is moved verbatim
apart from those visibility markers, four super::storage_api path
fixes, and the new per-module import headers.

Co-authored-by: heihutu <heihutu@gmail.com>

* refactor(scanner): split scanner_io.rs into child modules (#6294)

Split the 5369-line scanner_io.rs (39% inline tests) into a canonical
scanner_io.rs + scanner_io/ module tree with zero behavior change:

- scanner_io.rs (~660): constants, metadata-error constructors, the
  bucket scan plan, cycle-status classification helpers, the ScannerIO /
  ScannerIOCache / ScannerIODisk traits, and ScannerCycleResult
- scanner_io/dirty_usage.rs (~300): process-wide dirty-usage statics
  and the acknowledgment protocol
- scanner_io/guards.rs (~270): concurrency gauges and RAII guards
- scanner_io/cache.rs (~410): scanner cache locks and the snapshot
  persist/publish path
- scanner_io/io_cycle.rs (~390), io_cache.rs (~1160), io_disk.rs
  (~230): the ECStore / SetDisks / Disk trait implementations
- scanner_io/publish_gate_tests.rs (~750) and tests.rs (~1340): the two
  inline test modules as child modules

All crate paths are unchanged: the lib.rs scanner_io re-exports and
every crate::scanner_io:: consumer (scanner.rs, remote_scanner,
scanner_folder, and cross-crate rustfs users) resolve through root
re-exports with their original visibilities (pub stays pub, pub(crate)
stays pub(crate)). Cross-module items gain pub(super), whose scope
equals the old single-module privacy domain. Code is moved verbatim
apart from those markers, per-module import headers, and rustfmt
re-wraps.

The logging-guardrail nsscanner_disk skip-set_disks rule now points at
scanner_io/io_disk.rs where the function moved; the pattern and
thresholds are unchanged.

Co-authored-by: heihutu <heihutu@gmail.com>

* refactor(scanner): split data_usage_define persistence and tests (#6292)

Split the 3655-line data_usage_define.rs (59% inline tests) into a
canonical foo.rs + foo/ module tree with zero behavior change:

- data_usage_define.rs (~950): cache constants and revision helpers,
  the data-usage tree types, DataUsageCacheInfo with its hand-written
  Serialize, the in-memory tree operations, dui, and marshal/unmarshal
- data_usage_define/persistence.rs (~580): the load/backup/restore
  ladder (load, try_load_inner, revision_for_path) and the CAS save
  path with its retry policy and save metrics
- data_usage_define/tests.rs (~2155): the inline test module as a child
  module

All module paths are unchanged (the lib.rs data_usage_define::* glob
re-export and every crate::data_usage_define:: consumer resolve as
before). The hand-written map-encoded Serialize for
DataUsageCacheInfo is moved byte-for-byte per the AGENTS.md
cross-cutting invariant; on-disk names and the cache key format const
stay in the root. Four persistence helpers used by tests gain
pub(super), whose scope equals the old single-module privacy domain.
Code is moved verbatim apart from those markers, per-module import
headers, and rustfmt re-wraps.

Co-authored-by: heihutu <heihutu@gmail.com>

* chore(deps): bump datafusion to 55.0.0 (#6288)

* refactor(heal): split task.rs per heal kind (#6293)

* feat(ecstore): batch small file fdatasync commits (#6297)

* feat(ecstore): batch small file fdatasync commits

Add a default-off experimental file fdatasync group commit path for small rename_data shard directories. The coordinator batches same-disk waiters into one blocking task while preserving per-directory source fsync after shard contents are durable.

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(e2e): wait for compression S3 readiness

Reuse the shared S3 API readiness probe for compression test servers so multipart requests do not race the startup readiness gate after the TCP port opens.

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>

* fix(tier): recover multi-committed mutation intents (#6296)

* fix(tier): recover multi-committed mutation intents

* fix(tier): recover committed mutations on standalone nodes

* test(scanner): avoid stack overflow in transition test

---------

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: cxymds <cxymds@gmail.com>

---------

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
cxymds
2026-08-20 17:42:09 +08:00
committed by GitHub
parent 621fcb93c7
commit 319a03e638
15 changed files with 3251 additions and 529 deletions
+748 -14
View File
@@ -558,8 +558,11 @@ mod tests {
lifecycle::{TRANSITION_PENDING, TransitionOptions},
tier_delete_journal::{
TIER_DELETE_JOURNAL_PREFIX, persist_tier_delete_journal_entry, recover_tier_delete_journal_entries,
tier_delete_journal_object_name,
},
tier_sweeper::{
Jentry, TierDeleteJournalState, TierDeleteSourceIdentity, transitioned_delete_journal_entry_for_source,
},
tier_sweeper::Jentry,
transition_transaction::{
TRANSITION_TRANSACTION_RECORD_PREFIX, TransitionCleanupDecision, TransitionCleanupProof, TransitionOperatorError,
TransitionOperatorProbe, TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode,
@@ -569,9 +572,10 @@ mod tests {
recover_transition_transaction_records, save_transition_transaction_record,
},
},
bucket::metadata::{BUCKET_LIFECYCLE_CONFIG, BUCKET_VERSIONING_CONFIG},
client::transition_api::ReaderImpl,
config::com,
disk::RUSTFS_META_BUCKET,
disk::{RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE},
runtime::{global::set_object_store_resolver, sources as runtime_sources},
services::tier::{
test_util::{MockWarmBackend, MockWarmOp, TransitionCleanupStoreBarrier, register_mock_tier},
@@ -608,6 +612,8 @@ mod tests {
use rustfs_config::server_config::KVS;
use rustfs_filemeta::ObjectPartInfo;
#[cfg(feature = "test-util")]
use rustfs_filemeta::{FileInfo, FileMeta};
#[cfg(feature = "test-util")]
use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase};
use rustfs_rio::{Checksum, ChecksumType};
use rustfs_utils::{
@@ -1114,21 +1120,41 @@ mod tests {
Arc<crate::store::ECStore>,
CancellationToken,
) {
let mut pools = Vec::with_capacity(pool_drive_counts.len());
for (pool_index, &drives_per_set) in pool_drive_counts.iter().enumerate() {
let mut endpoints = Vec::with_capacity(drives_per_set);
for disk_index in 0..drives_per_set {
let path = temp_dir.join(format!("pool{pool_index}/disk{disk_index}"));
tokio::fs::create_dir_all(&path).await.expect("create disk dir");
let mut endpoint = Endpoint::try_from(path.to_str().expect("disk path should be utf-8")).expect("local endpoint");
endpoint.set_pool_index(pool_index);
endpoint.set_set_index(0);
endpoint.set_disk_index(disk_index);
endpoints.push(endpoint);
let pool_layouts = pool_drive_counts
.iter()
.map(|&drives_per_set| (1, drives_per_set))
.collect::<Vec<_>>();
build_isolated_test_store_with_layout(temp_dir, cmd_line, &pool_layouts, shutdown).await
}
async fn build_isolated_test_store_with_layout(
temp_dir: &std::path::Path,
cmd_line: &str,
pool_layouts: &[(usize, usize)],
shutdown: CancellationToken,
) -> (
Arc<crate::runtime::instance::InstanceContext>,
Arc<crate::store::ECStore>,
CancellationToken,
) {
let mut pools = Vec::with_capacity(pool_layouts.len());
for (pool_index, &(set_count, drives_per_set)) in pool_layouts.iter().enumerate() {
let mut endpoints = Vec::with_capacity(set_count * drives_per_set);
for set_index in 0..set_count {
for disk_index in 0..drives_per_set {
let path = temp_dir.join(format!("pool{pool_index}/set{set_index}/disk{disk_index}"));
tokio::fs::create_dir_all(&path).await.expect("create disk dir");
let mut endpoint =
Endpoint::try_from(path.to_str().expect("disk path should be utf-8")).expect("local endpoint");
endpoint.set_pool_index(pool_index);
endpoint.set_set_index(set_index);
endpoint.set_disk_index(disk_index);
endpoints.push(endpoint);
}
}
pools.push(PoolEndpoints {
legacy: false,
set_count: 1,
set_count,
drives_per_set,
endpoints: Endpoints::from(endpoints),
cmd_line: format!("{cmd_line}-pool-{pool_index}"),
@@ -2726,6 +2752,67 @@ mod tests {
.expect_err("suspended delete must remove the requested UUID version");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn tiered_data_movement_rejects_a_stale_source_snapshot_before_target_write() {
let temp_dir = tempfile::tempdir().expect("create stale-source data movement store dir");
let (_ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "stale-tier-source", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = "stale-tier-source-bucket";
let object = "stale-tier-source-object";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let incarnation = store
.bucket_incarnation_id(bucket)
.await
.expect("bucket incarnation should exist");
let version_id = uuid::Uuid::new_v4();
let stale = FileInfo {
volume: bucket.to_string(),
name: object.to_string(),
version_id: Some(version_id),
transition_status: rustfs_filemeta::TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/stale-source".to_string(),
transition_tier: "STALE-TIER".to_string(),
transition_version_id: Some(uuid::Uuid::new_v4()),
transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
data_dir: Some(uuid::Uuid::new_v4()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
size: 1,
metadata: HashMap::new(),
..Default::default()
};
let err = store
.decommission_tiered_object(
bucket,
object,
&stale,
&ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
src_pool_idx: 0,
data_movement: true,
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
},
)
.await
.expect_err("a source removed after queue capture must fail closed");
assert!(matches!(err, Error::ObjectNotFound(_, _) | Error::FileNotFound));
let target_err = store.pools[1]
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect_err("stale source rejection must not write target metadata");
assert!(matches!(
target_err,
StorageError::ObjectNotFound(_, _) | StorageError::VersionNotFound(_, _, _)
));
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn data_movement_put_conflict_validates_only_selected_target_pool() {
@@ -4468,6 +4555,653 @@ mod tests {
shutdown_b.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn prepared_tier_delete_recovery_finds_directory_source_on_encoded_set() {
let temp_dir = tempfile::tempdir().expect("create temp store dir");
let shutdown = CancellationToken::new();
let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store_with_layout(
temp_dir.path(),
"prepared-directory-recovery",
&[(2, 4)],
shutdown,
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let pool = &store.pools[0];
let object = (0..10_000)
.map(|index| format!("directory-{index}/"))
.find(|candidate| {
let encoded = rustfs_utils::path::encode_dir_object(candidate);
!Arc::ptr_eq(&pool.get_disks_by_key(candidate), &pool.get_disks_by_key(&encoded))
})
.expect("test topology should have a directory key whose encoded form hashes to another set");
let encoded = rustfs_utils::path::encode_dir_object(&object);
assert!(!Arc::ptr_eq(&pool.get_disks_by_key(&object), &pool.get_disks_by_key(&encoded)));
let tier_name = "PREPAREDDIRECTORY";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve")
.backend_identity();
let bucket = "prepared-directory-recovery-bucket";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(b"directory source".to_vec());
let original = store
.put_object(bucket, &object, &mut reader, &ObjectOptions::default())
.await
.expect("directory source should be written");
store
.transition_object(
bucket,
&object,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: original.etag.clone().expect("source should have an ETag"),
..Default::default()
},
mod_time: original.mod_time,
..Default::default()
},
)
.await
.expect("directory transition should commit");
let committed = store
.get_object_info(
bucket,
&object,
&ObjectOptions {
no_lock: true,
metadata_cache_safe: false,
..Default::default()
},
)
.await
.expect("transitioned directory source should be readable");
let mut entry = transitioned_delete_journal_entry_for_source(None, false, false, bucket, &object, &committed)
.expect("transitioned source should produce a prepared journal entry");
entry.backend_identity = Some(backend_identity);
persist_tier_delete_journal_entry(store.clone(), &entry)
.await
.expect("prepared journal should persist");
let stats = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("prepared recovery should complete");
assert_eq!((stats.scanned, stats.deleted, stats.failed), (1, 1, 0));
assert_eq!(tier_delete_journal_count(store).await, 0);
assert_eq!(backend.remove_count().await, 0);
assert_eq!(backend.object_count().await, 1, "live directory source must retain its remote object");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn prepared_tier_delete_recovery_checks_later_pool_then_commits_after_source_removal() {
let temp_dir = tempfile::tempdir().expect("create cross-pool recovery store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "prepared-cross-pool-recovery", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "PREPAREDCROSSPOOL";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve")
.backend_identity();
let bucket = "prepared-cross-pool-recovery-bucket";
let object = "object";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(b"later pool source".to_vec());
let original = store.pools[1]
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("source should be written only to the later pool");
store.pools[1]
.transition_object(
bucket,
object,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: original.etag.clone().expect("source should have an ETag"),
..Default::default()
},
mod_time: original.mod_time,
..Default::default()
},
)
.await
.expect("later-pool transition should commit");
let committed = store.pools[1]
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("transitioned source should be readable");
let mut entry = transitioned_delete_journal_entry_for_source(None, false, false, bucket, object, &committed)
.expect("transitioned source should produce a prepared journal");
entry.backend_identity = Some(backend_identity);
persist_tier_delete_journal_entry(store.clone(), &entry)
.await
.expect("prepared journal should persist");
let retained = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("recovery should scan the later pool");
assert_eq!((retained.scanned, retained.deleted, retained.failed), (1, 1, 0));
assert_eq!(backend.remove_count().await, 0);
assert_eq!(backend.object_count().await, 1);
assert_eq!(
tier_delete_journal_count(store.clone()).await,
0,
"live source should abort its prepared journal"
);
store.pools[1]
.delete_object(
bucket,
object,
ObjectOptions {
version_id: committed.version_id.map(|version| version.to_string()),
expiration: crate::storage_api_contracts::lifecycle::ExpirationOptions { expire: true },
..Default::default()
},
)
.await
.expect("source version should be removed before recovery retry");
persist_tier_delete_journal_entry(store.clone(), &entry)
.await
.expect("prepared journal should persist for the absent source");
let deleted = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("recovery should commit an absent stable source");
assert_eq!((deleted.scanned, deleted.deleted, deleted.failed), (1, 1, 0));
assert_eq!(tier_delete_journal_count(store).await, 0);
assert_eq!(backend.remove_count().await, 1);
assert_eq!(backend.object_count().await, 0);
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn prepared_tier_delete_recovery_retains_journal_on_source_metadata_error() {
let temp_dir = tempfile::tempdir().expect("create metadata-error recovery store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "prepared-metadata-error-recovery", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "PREPAREDMETADATAERROR";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve")
.backend_identity();
let bucket = "prepared-metadata-error-recovery-bucket";
let object = "object";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(b"source with unreadable metadata".to_vec());
let original = store
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("source should be written");
store
.transition_object(
bucket,
object,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: original.etag.clone().expect("source should have an ETag"),
..Default::default()
},
mod_time: original.mod_time,
..Default::default()
},
)
.await
.expect("source transition should commit");
let committed = store
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("transitioned source should be readable");
let mut entry = transitioned_delete_journal_entry_for_source(None, false, false, bucket, object, &committed)
.expect("transitioned source should produce a prepared journal");
entry.backend_identity = Some(backend_identity);
persist_tier_delete_journal_entry(store.clone(), &entry)
.await
.expect("prepared journal should persist");
for disk_index in 0..4 {
let metadata_path = temp_dir
.path()
.join(format!("pool0/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
tokio::fs::write(metadata_path, b"not-xl-meta")
.await
.expect("source metadata should be corrupted");
}
let stats = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("recovery scan should complete despite the entry failure");
assert_eq!((stats.scanned, stats.deleted, stats.failed), (1, 0, 1));
assert_eq!(tier_delete_journal_count(store).await, 1, "journal must remain prepared for retry");
assert_eq!(backend.remove_count().await, 0, "unreadable source metadata must block remote deletion");
assert_eq!(backend.object_count().await, 1, "remote source must remain intact");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn tier_delete_recovery_retains_content_that_does_not_match_its_object_name() {
let temp_dir = tempfile::tempdir().expect("create mismatched journal store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "mismatched-tier-journal", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "MISMATCHEDJOURNAL";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve");
let remote_version = uuid::Uuid::new_v4().to_string();
backend.set_put_remote_version(Some(remote_version.clone())).await;
lease
.put("remote/original", ReaderImpl::Body(bytes::Bytes::from_static(b"remote body")), 11)
.await
.expect("remote body should be seeded");
let entry = Jentry {
obj_name: "remote/original".to_string(),
version_id: remote_version,
tier_name: tier_name.to_string(),
backend_identity: Some(lease.backend_identity()),
version_id_exact: true,
version_state: rustfs_filemeta::TransitionVersionState::Exact,
state: TierDeleteJournalState::Prepared,
source: Some(TierDeleteSourceIdentity {
bucket: "absent-source-bucket".to_string(),
object: "absent-source-object".to_string(),
version_id: Some(uuid::Uuid::new_v4().to_string()),
versioned: true,
version_suspended: false,
data_dir: Some(uuid::Uuid::new_v4().to_string()),
etag: Some("etag".to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH.to_string()),
}),
};
persist_tier_delete_journal_entry(store.clone(), &entry)
.await
.expect("prepared journal should persist");
let journal_name = tier_delete_journal_object_name(&entry);
let data = com::read_config(store.clone(), &journal_name)
.await
.expect("prepared journal should be readable");
let mut value: serde_json::Value = serde_json::from_slice(&data).expect("journal should contain JSON");
value["obj_name"] = serde_json::json!("remote/replaced");
com::save_config(
store.clone(),
&journal_name,
serde_json::to_vec(&value).expect("mismatched journal should encode"),
)
.await
.expect("mismatched journal content should be written under the original name");
let stats = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("recovery scan should complete");
assert_eq!((stats.scanned, stats.deleted, stats.failed), (1, 0, 1));
assert_eq!(tier_delete_journal_count(store).await, 1, "mismatched journal must be retained");
assert_eq!(backend.remove_count().await, 0);
assert_eq!(backend.object_count().await, 1);
}
#[cfg(feature = "test-util")]
#[test]
#[serial_test::serial(storage_class_env)]
fn transitioned_history_expiry_journals_real_source_without_free_version() {
std::thread::Builder::new()
.name("transitioned-delete-all-test".to_string())
.stack_size(32 * 1024 * 1024)
.spawn(|| {
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.worker_threads(2)
.build()
.expect("test runtime should build");
runtime.block_on(async {
let temp_dir = tempfile::tempdir().expect("create transitioned delete-all store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "transitioned-delete-all", &[4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "DELETEALLTRANSITIONED";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let bucket = "transitioned-delete-all-bucket";
let object = "object";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
crate::bucket::metadata_sys::update(
bucket,
BUCKET_VERSIONING_CONFIG,
b"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>".to_vec(),
)
.await
.expect("bucket versioning should be enabled");
crate::bucket::metadata_sys::update(
bucket,
BUCKET_LIFECYCLE_CONFIG,
br#"<LifecycleConfiguration>
<Rule>
<ID>delete-all-versions</ID>
<Status>Enabled</Status>
<Filter><Prefix></Prefix></Filter>
<Expiration><Days>1</Days><ExpiredObjectAllVersions>true</ExpiredObjectAllVersions></Expiration>
</Rule>
</LifecycleConfiguration>"#
.to_vec(),
)
.await
.expect("delete-all lifecycle should be configured");
let old_time = OffsetDateTime::now_utc() - time::Duration::days(3);
let mut history_reader = PutObjReader::from_vec(b"transitioned history".to_vec());
let history = store
.put_object(
bucket,
object,
&mut history_reader,
&ObjectOptions {
versioned: true,
mod_time: Some(old_time - time::Duration::hours(1)),
..Default::default()
},
)
.await
.expect("historical version should be written");
let mut current_reader = PutObjReader::from_vec(b"current version".to_vec());
let current = store
.put_object(
bucket,
object,
&mut current_reader,
&ObjectOptions {
versioned: true,
mod_time: Some(old_time),
..Default::default()
},
)
.await
.expect("current version should be written");
store
.transition_object(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: history.version_id.map(|version_id| version_id.to_string()),
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: history.etag.clone().expect("history should have an ETag"),
..Default::default()
},
mod_time: history.mod_time,
..Default::default()
},
)
.await
.expect("historical version should transition");
assert_eq!(backend.object_count().await, 1);
let transitioned_remote_versions = backend.put_versions().await;
assert_eq!(transitioned_remote_versions.len(), 1);
let incarnation = store
.bucket_incarnation_id_from_disk(bucket)
.await
.expect("bucket incarnation should be available");
for disk_index in 0..4 {
let metadata_path = temp_dir
.path()
.join(format!("pool0/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let encoded = tokio::fs::read(&metadata_path)
.await
.expect("transition metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("transition metadata should decode");
let mut transitioned = metadata
.get_all_file_info_versions(bucket, object, true)
.expect("transitioned versions should decode")
.versions
.into_iter()
.find(|version| version.version_id == history.version_id)
.expect("transitioned history should exist");
transitioned.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown;
metadata
.add_version(transitioned)
.expect("unknown state should replace the transitioned version");
tokio::fs::write(
&metadata_path,
metadata.marshal_msg().expect("unknown transition metadata should encode"),
)
.await
.expect("unknown transition metadata should be written");
}
let lifecycle_event = crate::bucket::lifecycle::lifecycle::Event {
action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "delete-all-versions".to_string(),
..Default::default()
};
let rejected = crate::bucket::lifecycle::bucket_lifecycle_ops::apply_expiry_on_non_transitioned_objects(
store.clone(),
&current,
&lifecycle_event,
&crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc::Scanner,
incarnation,
)
.await;
assert!(!rejected, "legacy unknown transition identity must fail before local mutation");
assert_eq!(tier_delete_journal_count(store.clone()).await, 0);
assert_eq!(backend.remove_count().await, 0);
let retained = store.pools[0].disk_set[0]
.load_file_info_versions_exact(bucket, object)
.await
.expect("rejected delete-all metadata should remain readable")
.expect("rejected delete-all should retain both versions");
assert_eq!(
retained
.versions
.iter()
.filter(|version| !version.tier_free_version())
.count(),
2
);
for disk_index in 0..4 {
let metadata_path = temp_dir
.path()
.join(format!("pool0/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let encoded = tokio::fs::read(&metadata_path)
.await
.expect("unknown transition metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("unknown transition metadata should decode");
let mut transitioned = metadata
.get_all_file_info_versions(bucket, object, true)
.expect("unknown transition versions should decode")
.versions
.into_iter()
.find(|version| version.version_id == history.version_id)
.expect("unknown transitioned history should exist");
transitioned.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
metadata
.add_version(transitioned)
.expect("exact state should replace the transitioned version");
tokio::fs::write(
&metadata_path,
metadata.marshal_msg().expect("exact transition metadata should encode"),
)
.await
.expect("exact transition metadata should be written");
}
let applied = crate::bucket::lifecycle::bucket_lifecycle_ops::apply_expiry_on_non_transitioned_objects(
store.clone(),
&current,
&lifecycle_event,
&crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc::Scanner,
incarnation,
)
.await;
assert!(applied, "delete-all should remove current and transitioned history");
let versions = store.pools[0].disk_set[0]
.load_file_info_versions_exact(bucket, object)
.await
.expect("remaining exact metadata should be readable");
assert!(versions.is_none(), "delete-all must not leave a tier free-version");
assert_eq!(tier_delete_journal_count(store.clone()).await, 1);
assert_eq!(backend.object_count().await, 1, "remote deletion must remain journal-driven");
let stats = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("committed journal should recover");
assert_eq!((stats.scanned, stats.deleted, stats.failed), (1, 1, 0));
assert_eq!(tier_delete_journal_count(store).await, 0);
assert_eq!(backend.object_count().await, 0);
assert_eq!(backend.exact_remove_count(), 1);
assert_eq!(backend.remove_versions().await, transitioned_remote_versions);
});
})
.expect("test thread should spawn")
.join()
.expect("test thread should complete");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn prepared_tier_delete_recovery_requires_namespace_locking() {
temp_env::async_with_vars([("RUSTFS_LOCK_ENABLED", Some("false"))], async {
let temp_dir = tempfile::tempdir().expect("create lock-disabled store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "prepared-recovery-lock-disabled", &[4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
assert!(ctx.lock_manager().is_disabled());
let tier_name = "PREPAREDLOCKDISABLED";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("tier lease should resolve")
.backend_identity();
let entry = Jentry {
obj_name: "remote/lock-disabled".to_string(),
version_id: uuid::Uuid::new_v4().to_string(),
tier_name: tier_name.to_string(),
backend_identity: Some(backend_identity),
version_id_exact: true,
version_state: rustfs_filemeta::TransitionVersionState::Exact,
state: TierDeleteJournalState::Prepared,
source: Some(TierDeleteSourceIdentity {
bucket: "absent-source-bucket".to_string(),
object: "absent-source-object".to_string(),
version_id: Some(uuid::Uuid::new_v4().to_string()),
versioned: true,
version_suspended: false,
data_dir: Some(uuid::Uuid::new_v4().to_string()),
etag: Some("etag".to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH.to_string()),
}),
};
persist_tier_delete_journal_entry(store.clone(), &entry)
.await
.expect("prepared journal should persist");
let stats = recover_tier_delete_journal_entries(store.clone(), 100, None)
.await
.expect("recovery scan should complete");
assert_eq!((stats.scanned, stats.deleted, stats.failed), (1, 0, 1));
assert_eq!(tier_delete_journal_count(store).await, 1, "journal must remain prepared for retry");
assert_eq!(backend.remove_count().await, 0, "lock-disabled recovery must not delete remotely");
})
.await;
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn lifecycle_delete_all_requires_namespace_locking_before_mutation() {
temp_env::async_with_vars([("RUSTFS_LOCK_ENABLED", Some("false"))], async {
let temp_dir = tempfile::tempdir().expect("create lock-disabled delete-all store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "delete-all-lock-disabled", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
assert!(ctx.lock_manager().is_disabled());
let bucket = "delete-all-lock-disabled-bucket";
let object = "object";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(b"must survive".to_vec());
let no_lock_opts = ObjectOptions {
no_lock: true,
..Default::default()
};
let original = store.pools[0]
.put_object(bucket, object, &mut reader, &no_lock_opts)
.await
.expect("source should be written");
let mut delete_opts = ObjectOptions {
delete_prefix: true,
delete_prefix_object: true,
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: original.version_id,
delete_marker: false,
action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
delete_replication_config_snapshot: Some(Arc::new(
crate::bucket::replication::DeleteReplicationConfigSnapshot::default(),
)),
..Default::default()
};
delete_opts.ensure_lifecycle_delete_all_journal();
let err = store
.delete_object_with_tier_delete_journal(bucket, object, delete_opts)
.await
.expect_err("delete-all must reject disabled namespace locking");
assert!(err.to_string().contains("requires namespace locking"));
let retained = store.pools[0]
.get_object_info(bucket, object, &no_lock_opts)
.await
.expect("rejected delete-all must retain the source");
assert_eq!(retained.etag, original.etag);
assert_eq!(tier_delete_journal_count(store).await, 0, "rejected delete-all must not prepare journals");
})
.await;
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
+1
View File
@@ -151,6 +151,7 @@ pub(crate) mod init_format;
pub(crate) mod list_objects;
mod multipart;
mod object;
pub(crate) use object::ObjectLockDiagGuard;
pub use object::{
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
SnapshotConsistencyError,
+213 -22
View File
@@ -14,6 +14,8 @@
use super::*;
use crate::bucket::lifecycle::{
bucket_lifecycle_ops::eval_action_from_lifecycle,
get_expiry_configs,
tier_delete_journal::{
abort_prepared_tier_delete_journal_entry as abort_prepared_journal_entry_if_current, commit_tier_delete_journal_entry,
enqueue_committed_tier_delete_journal_entry, persist_tier_delete_journal_entry,
@@ -211,7 +213,8 @@ async fn delete_prefix_with_tier_delete_journal(
opts: &ObjectOptions,
tier_journal_api: Option<&Arc<ECStore>>,
) -> Result<()> {
let journal_entry = if let Some(api) = tier_journal_api {
let lifecycle_delete_all = opts.lifecycle_delete_all.is_some();
let journal_entry = if !lifecycle_delete_all && let Some(api) = tier_journal_api {
Some(prepare_prefix_tier_delete_journal_entries(api, bucket, object, opts).await?)
} else {
None
@@ -220,14 +223,34 @@ async fn delete_prefix_with_tier_delete_journal(
let result = store.delete_prefix(bucket, object, opts).await;
match result {
Ok(()) => {
if let (Some(api), Some(entries)) = (tier_journal_api, journal_entry.as_ref()) {
let lifecycle_entries = if lifecycle_delete_all {
opts.lifecycle_delete_all_journal()
.ok_or(StorageError::PreconditionFailed)?
.lock()
.prepared_entries()
} else {
Vec::new()
};
let entries = journal_entry.as_deref().unwrap_or(&lifecycle_entries);
if let Some(api) = tier_journal_api {
commit_prepared_tier_delete_journal_entries(api, entries).await;
}
Ok(())
}
Err(err) => {
if let (Some(api), Some(entries)) = (tier_journal_api, journal_entry.as_ref()) {
abort_prepared_tier_delete_journal_entries(api, entries).await;
if let Some(api) = tier_journal_api {
if lifecycle_delete_all {
let (abort, entries) = {
let journal = opts.lifecycle_delete_all_journal().ok_or(StorageError::PreconditionFailed)?;
let state = journal.lock();
(!state.mutation_started(), state.prepared_entries())
};
if abort {
abort_prepared_tier_delete_journal_entries(api, &entries).await;
}
} else if let Some(entries) = journal_entry.as_ref() {
abort_prepared_tier_delete_journal_entries(api, entries).await;
}
}
Err(err)
}
@@ -327,7 +350,7 @@ impl fmt::Display for ObjectLockDiagMode {
}
}
struct ObjectLockDiagGuard {
pub(crate) struct ObjectLockDiagGuard {
guard: rustfs_lock::NamespaceLockGuard,
enabled: bool,
op: &'static str,
@@ -360,14 +383,14 @@ impl ObjectLockDiagGuard {
}
}
fn lock_lost_signal(&self) -> Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>> {
pub(crate) fn lock_lost_signal(&self) -> Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>> {
match &self.guard {
rustfs_lock::NamespaceLockGuard::Standard(guard) => Some(guard.lock_lost()),
rustfs_lock::NamespaceLockGuard::Fast(_) => None,
}
}
fn is_lock_lost(&self) -> bool {
pub(crate) fn is_lock_lost(&self) -> bool {
self.guard.is_lock_lost()
}
}
@@ -1109,6 +1132,26 @@ fn is_equivalent_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo,
&& source_actual_size == target_actual_size
}
fn tiered_data_movement_source_matches(
expected: &rustfs_filemeta::FileInfo,
current: &rustfs_filemeta::FileInfo,
) -> Result<bool> {
let expected_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&expected.metadata)?;
let current_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&current.metadata)?;
Ok(expected.version_id == current.version_id
&& expected.data_dir == current.data_dir
&& expected.mod_time == current.mod_time
&& expected.size == current.size
&& expected.get_etag() == current.get_etag()
&& expected.transition_status == current.transition_status
&& expected.transitioned_objname == current.transitioned_objname
&& expected.transition_tier == current.transition_tier
&& expected.transition_version_id == current.transition_version_id
&& expected.transition_version == current.transition_version
&& expected.transition_version_state == current.transition_version_state
&& expected_backend == current_backend)
}
fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool {
target_pool_idx != src_pool_idx
}
@@ -1247,7 +1290,9 @@ impl ECStore {
let mut opts = opts.clone();
opts.no_lock = false;
opts.metadata_cache_safe = false;
let read_lock_guards = self.acquire_select_object_read_locks(bucket, &object, &mut opts).await?;
let read_lock_guards = self
.acquire_all_object_read_locks("select_object", bucket, &object, &mut opts)
.await?;
if self.ctx.lock_manager().is_disabled() {
return Err(SnapshotConsistencyError::LockingDisabled.into());
}
@@ -1488,8 +1533,9 @@ impl ECStore {
)))
}
async fn acquire_select_object_read_locks(
pub(crate) async fn acquire_all_object_read_locks(
&self,
op: &'static str,
bucket: &str,
object: &str,
opts: &mut ObjectOptions,
@@ -1501,10 +1547,7 @@ impl ECStore {
// for each object's hashed set. DELETE and same-key CopyObject use the
// fixed domain, while PUT commits and data movement use the hashed set.
let distributed = self.ctx.is_dist_erasure().await;
if let Some(guard) = self
.acquire_object_read_lock_if_needed("select_object", bucket, object, opts)
.await?
{
if let Some(guard) = self.acquire_object_read_lock_if_needed(op, bucket, object, opts).await? {
guards.push(guard);
}
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
@@ -1527,7 +1570,7 @@ impl ECStore {
.map_err(|err| Self::map_namespace_lock_error(bucket, object, "read", err))?;
let owner = diag_enabled.then(|| ns_lock.owner().to_string());
log_object_lock_acquire_if_slow(
"select_object",
op,
bucket,
object,
owner.as_deref(),
@@ -1538,7 +1581,7 @@ impl ECStore {
guards.push(ObjectLockDiagGuard::new(
guard,
diag_enabled,
"select_object",
op,
diag_enabled.then(|| bucket.to_string()),
diag_enabled.then(|| object.to_string()),
owner,
@@ -1549,6 +1592,77 @@ impl ECStore {
Ok(guards)
}
async fn acquire_data_movement_object_write_locks(
&self,
bucket: &str,
object: &str,
source_pool_idx: usize,
target_pool_idx: usize,
opts: &mut ObjectOptions,
) -> Result<Vec<ObjectLockDiagGuard>> {
if self.ctx.lock_manager().is_disabled() {
return Err(Error::other("tiered data movement requires namespace locking"));
}
let distributed = self.ctx.is_dist_erasure().await;
let diag_enabled = is_object_lock_diag_enabled();
let mut pool_indices = [source_pool_idx, target_pool_idx];
pool_indices.sort_unstable();
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
let mut locked_sets = vec![fixed_set];
let mut guards = Vec::with_capacity(3);
// Lock order matches journal recovery: fixed store domain first, then
// hashed domains by ascending pool index. This also serializes source
// revalidation and target publication against ordinary object deletes.
guards.push(self.acquire_object_write_lock("tiered_data_movement", bucket, object).await?);
for pool_idx in pool_indices {
let pool = self
.pools
.get(pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement pool {pool_idx}")))?;
let set = pool.get_disks_by_key(object);
let lock_domain_already_held = !distributed
|| locked_sets.iter().any(|locked_set: &Arc<crate::set_disk::SetDisks>| {
same_distributed_lock_domain(&locked_set.lockers, &set.lockers)
});
if lock_domain_already_held {
continue;
}
let ns_lock = set.new_ns_lock(bucket, object).await?;
let acquire_start = Instant::now();
let guard = ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(|err| Self::map_namespace_lock_error(bucket, object, "write", err))?;
let owner = diag_enabled.then(|| ns_lock.owner().to_string());
log_object_lock_acquire_if_slow(
"tiered_data_movement",
bucket,
object,
owner.as_deref(),
ObjectLockDiagMode::Write,
acquire_start.elapsed(),
diag_enabled,
);
guards.push(ObjectLockDiagGuard::new(
guard,
diag_enabled,
"tiered_data_movement",
diag_enabled.then(|| bucket.to_string()),
diag_enabled.then(|| object.to_string()),
owner,
ObjectLockDiagMode::Write,
));
locked_sets.push(set);
}
opts.no_lock = true;
for signal in guards.iter().filter_map(ObjectLockDiagGuard::lock_lost_signal) {
opts.add_namespace_lock_lost_signal(signal);
}
opts.ensure_namespace_lock_fence();
Ok(guards)
}
fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option<ObjectLockDiagGuard>) -> GetObjectReader {
if is_lock_optimization_enabled() || reader.buffered_body.is_some() {
return reader;
@@ -1686,13 +1800,8 @@ impl ECStore {
Some(guard)
};
let mut fi = fi.clone();
if opts.data_movement {
crate::data_movement::prepare_tiered_data_movement_file_info(&mut fi)?;
}
let object = encode_dir_object(object);
let logical_object = object;
let object = encode_dir_object(logical_object);
if self.single_pool() {
return Self::resolve_decommission_tiered_object_result(
Err(Error::other("single pool deployments cannot decommission tiered objects")),
@@ -1715,6 +1824,33 @@ impl ECStore {
&object,
)?
};
let _object_guards = self
.acquire_data_movement_object_write_locks(bucket, &object, opts.src_pool_idx, idx, &mut opts)
.await?;
let source_pool = self
.pools
.get(opts.src_pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement source pool {}", opts.src_pool_idx)))?;
let source_versions = source_pool
.get_disks_by_key(&object)
.load_file_info_versions_exact(bucket, logical_object)
.await?;
let current_source = source_versions
.as_ref()
.and_then(|versions| {
versions
.versions
.iter()
.find(|current| current.version_id == fi.version_id && !current.tier_free_version())
})
.ok_or_else(|| to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]))?;
if !tiered_data_movement_source_matches(fi, current_source)? {
return Err(to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]));
}
let mut fi = current_source.clone();
if opts.data_movement {
crate::data_movement::prepare_tiered_data_movement_file_info(&mut fi)?;
}
if opts.data_movement && idx == opts.src_pool_idx {
let resume_target_pool_idx = self
.get_available_pool_idx_excluding(bucket, &object, fi.size, opts.src_pool_idx)
@@ -2190,6 +2326,10 @@ impl ECStore {
) -> Result<ObjectInfo> {
check_del_obj_args(bucket, object)?;
if opts.lifecycle_delete_all.is_some() && self.ctx.lock_manager().is_disabled() {
return Err(Error::other("lifecycle delete-all requires namespace locking"));
}
let _bucket_lifecycle_guard = if is_meta_bucketname(bucket) {
None
} else if opts.delete_prefix {
@@ -2204,6 +2344,11 @@ impl ECStore {
};
let object = object.as_str();
let mut opts = opts;
let delete_all_configs = if opts.lifecycle_delete_all.is_some() {
Some(get_expiry_configs(self, bucket).await?)
} else {
None
};
opts.tier_delete_journal_api = tier_journal_api.clone();
if let Some(guard) = _bucket_lifecycle_guard.as_ref() {
opts.add_bucket_lifecycle_lock_guard(guard);
@@ -2301,6 +2446,34 @@ impl ECStore {
} else {
None
};
if let Some(trigger) = opts.lifecycle_delete_all.as_ref() {
let configs = delete_all_configs.as_ref().ok_or(StorageError::PreconditionFailed)?;
let expected_bucket_incarnation_id = opts.expected_bucket_incarnation_id.ok_or(StorageError::PreconditionFailed)?;
if configs.table_bucket_enabled || configs.bucket_incarnation_id != expected_bucket_incarnation_id {
return Err(StorageError::PreconditionFailed);
}
let lifecycle = configs.lifecycle.as_ref().ok_or(StorageError::PreconditionFailed)?;
let (mut current, _) = self
.get_latest_object_info_with_idx(
bucket,
object,
&ObjectOptions {
no_lock: true,
metadata_cache_safe: false,
..Default::default()
},
)
.await?;
let current_version_id = current.version_id.filter(|version_id| !version_id.is_nil());
if current_version_id != trigger.version_id || current.delete_marker != trigger.delete_marker {
return Err(StorageError::PreconditionFailed);
}
current.name = decode_dir_object(&current.name);
let current_event = eval_action_from_lifecycle(lifecycle, configs.object_lock.as_deref(), &current).await;
if current_event.action != trigger.action || current_event.rule_id != trigger.rule_id {
return Err(StorageError::PreconditionFailed);
}
}
if opts.delete_prefix {
delete_prefix_with_tier_delete_journal(self, bucket, object, &opts, tier_journal_api.as_ref()).await?;
return Ok(ObjectInfo::default());
@@ -3828,6 +4001,24 @@ mod tests {
assert!(is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn tiered_data_movement_source_match_rejects_transition_identity_changes() {
let source = tiered_equivalence_source();
assert!(tiered_data_movement_source_matches(&source, &source).expect("matching source metadata should parse"));
let mut changed_remote = source.clone();
changed_remote.transitioned_objname = "remote/replaced".to_string();
assert!(!tiered_data_movement_source_matches(&source, &changed_remote).expect("changed remote metadata should parse"));
let mut changed_backend = source.clone();
rustfs_utils::http::metadata_compat::insert_str(
&mut changed_backend.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex([9; 32]),
);
assert!(!tiered_data_movement_source_matches(&source, &changed_backend).expect("backend metadata should parse"));
}
#[test]
fn equivalent_data_movement_tiered_object_uses_logical_compressed_and_encrypted_sizes() {
let mut compressed = tiered_equivalence_source();
+346 -8
View File
@@ -202,10 +202,76 @@ impl ECStore {
}
pub(super) async fn delete_prefix(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> {
if opts.lifecycle_delete_all.is_some() {
let mut preflight_opts = opts.clone();
preflight_opts
.lifecycle_delete_all
.as_mut()
.ok_or(StorageError::PreconditionFailed)?
.phase = crate::object_api::LifecycleDeleteAllPhase::Preflight;
for pool in &self.pools {
#[cfg(test)]
lifecycle_delete_all_test_failure(crate::object_api::LifecycleDeleteAllPhase::Preflight, pool.pool_idx)?;
pool.delete_object(bucket, object, preflight_opts.clone()).await?;
}
opts.lifecycle_delete_all_journal()
.ok_or(StorageError::PreconditionFailed)?
.lock()
.mark_mutation_started();
let mut non_trigger_opts = opts.clone();
non_trigger_opts
.lifecycle_delete_all
.as_mut()
.ok_or(StorageError::PreconditionFailed)?
.phase = crate::object_api::LifecycleDeleteAllPhase::History;
for pool in &self.pools {
#[cfg(test)]
lifecycle_delete_all_test_failure(crate::object_api::LifecycleDeleteAllPhase::History, pool.pool_idx)?;
let mut pool_opts = non_trigger_opts.clone();
pool_opts.delete_prefix = true;
pool.delete_object(bucket, object, pool_opts).await?;
}
let mut final_preflight_opts = opts.clone();
final_preflight_opts
.lifecycle_delete_all
.as_mut()
.ok_or(StorageError::PreconditionFailed)?
.phase = crate::object_api::LifecycleDeleteAllPhase::FinalPreflight;
let mut trigger_pools = Vec::new();
for (pool_index, pool) in self.pools.iter().enumerate() {
#[cfg(test)]
lifecycle_delete_all_test_failure(crate::object_api::LifecycleDeleteAllPhase::FinalPreflight, pool.pool_idx)?;
let result = pool.delete_object(bucket, object, final_preflight_opts.clone()).await?;
if !result.name.is_empty() {
trigger_pools.push(pool_index);
}
}
if trigger_pools.is_empty() {
return Err(StorageError::PreconditionFailed);
}
let mut trigger_opts = opts.clone();
trigger_opts
.lifecycle_delete_all
.as_mut()
.ok_or(StorageError::PreconditionFailed)?
.phase = crate::object_api::LifecycleDeleteAllPhase::Trigger;
for pool_index in trigger_pools {
#[cfg(test)]
lifecycle_delete_all_test_failure(crate::object_api::LifecycleDeleteAllPhase::Trigger, pool_index)?;
let mut pool_opts = trigger_opts.clone();
pool_opts.delete_prefix = true;
self.pools[pool_index].delete_object(bucket, object, pool_opts).await?;
}
return Ok(());
}
let mut first_error = None;
let mut first_volume_error = None;
let mut has_success = false;
for pool in self.pools.iter() {
for pool in &self.pools {
let mut opts = opts.clone();
opts.delete_prefix = true;
match pool.delete_object(bucket, object, opts).await {
@@ -774,6 +840,22 @@ impl ECStore {
}
}
#[cfg(test)]
static LIFECYCLE_DELETE_ALL_TEST_FAILURE: std::sync::Mutex<Option<(crate::object_api::LifecycleDeleteAllPhase, usize)>> =
std::sync::Mutex::new(None);
#[cfg(test)]
fn lifecycle_delete_all_test_failure(phase: crate::object_api::LifecycleDeleteAllPhase, pool_index: usize) -> Result<()> {
if LIFECYCLE_DELETE_ALL_TEST_FAILURE
.lock()
.expect("lifecycle delete-all failure hook should not poison")
.is_some_and(|failure| failure == (phase, pool_index))
{
return Err(StorageError::PreconditionFailed);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
@@ -781,23 +863,28 @@ mod tests {
use crate::disk::error::DiskError;
use crate::layout::endpoint::Endpoint;
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::object_api::ObjectLockConfigSnapshot;
use crate::storage_api_contracts::bucket::MakeBucketOptions;
use crate::storage_api_contracts::object::ObjectIO as _;
use arc_swap::ArcSwap;
use rustfs_config::server_config::KVS;
use rustfs_filemeta::FileInfo;
use std::sync::Arc;
use tokio_util::sync::CancellationToken;
#[tokio::test]
async fn delete_prefix_attempts_later_pools_after_an_earlier_pool_error() {
let temp_dir = tempfile::tempdir().expect("multi-pool delete test directory should be created");
let mut pools = Vec::with_capacity(2);
for (pool_index, drives_per_set) in [2, 4].into_iter().enumerate() {
async fn setup_multi_pool_test_store(
name: &str,
drives_per_pool: &[usize],
) -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
let temp_dir = tempfile::tempdir().expect("multi-pool test directory should be created");
let mut pools = Vec::with_capacity(drives_per_pool.len());
for (pool_index, drives_per_set) in drives_per_pool.iter().copied().enumerate() {
let mut endpoints = Vec::with_capacity(drives_per_set);
for disk_index in 0..drives_per_set {
let disk_path = temp_dir.path().join(format!("pool{pool_index}-disk{disk_index}"));
tokio::fs::create_dir_all(&disk_path)
.await
.expect("multi-pool delete test disk should be created");
.expect("multi-pool test disk should be created");
let mut endpoint =
Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse");
endpoint.set_pool_index(pool_index);
@@ -810,7 +897,7 @@ mod tests {
set_count: 1,
drives_per_set,
endpoints: Endpoints::from(endpoints),
cmd_line: format!("delete-prefix-pool-{pool_index}"),
cmd_line: format!("{name}-pool-{pool_index}"),
platform: "test".to_string(),
});
}
@@ -830,6 +917,92 @@ mod tests {
.await
.expect("multi-pool store should initialize");
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
(temp_dir, store, shutdown)
}
struct LifecycleDeleteAllFailureGuard;
impl Drop for LifecycleDeleteAllFailureGuard {
fn drop(&mut self) {
*LIFECYCLE_DELETE_ALL_TEST_FAILURE
.lock()
.expect("lifecycle delete-all failure hook should not poison") = None;
}
}
async fn seed_multi_pool_delete_all(store: &Arc<ECStore>, bucket: &str, object: &str) -> ObjectOptions {
let trigger_id = Uuid::new_v4();
for (pool_index, pool) in store.pools.iter().enumerate() {
let mut history_reader = PutObjReader::from_vec(format!("{object}-history-{pool_index}").into_bytes());
pool.put_object(
bucket,
object,
&mut history_reader,
&ObjectOptions {
versioned: true,
version_id: Some(Uuid::new_v4().to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(1)),
..Default::default()
},
)
.await
.expect("history should be stored");
let mut trigger_reader = PutObjReader::from_vec(format!("{object}-trigger-{pool_index}").into_bytes());
pool.put_object(
bucket,
object,
&mut trigger_reader,
&ObjectOptions {
versioned: true,
version_id: Some(trigger_id.to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(2)),
..Default::default()
},
)
.await
.expect("shared trigger should be stored");
}
let mut opts = ObjectOptions {
delete_prefix: true,
delete_prefix_object: true,
versioned: true,
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(trigger_id),
delete_marker: false,
action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(
crate::bucket::metadata_sys::ObjectLockConfigState::ConfirmedAbsent,
))),
delete_replication_config_snapshot: Some(Arc::new(
crate::bucket::replication::DeleteReplicationConfigSnapshot::default(),
)),
..Default::default()
};
opts.ensure_lifecycle_delete_all_journal();
opts
}
async fn ordinary_version_count(store: &ECStore, pool_index: usize, bucket: &str, object: &str) -> usize {
store.pools[pool_index].disk_set[0]
.load_file_info_versions_exact(bucket, object)
.await
.expect("pool metadata should load")
.map(|versions| {
versions
.versions
.iter()
.filter(|version| !version.tier_free_version())
.count()
})
.unwrap_or_default()
}
#[tokio::test]
async fn delete_prefix_attempts_later_pools_after_an_earlier_pool_error() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("delete-prefix", &[2, 4]).await;
let bucket = format!("delete-prefix-{}", Uuid::new_v4().simple());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
@@ -921,6 +1094,171 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn lifecycle_delete_all_history_failure_preserves_trigger_and_retry_converges() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("lifecycle-delete-all", &[4, 4]).await;
let bucket = format!("lifecycle-delete-all-{}", Uuid::new_v4().simple());
let object = "object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created in both pools");
for pool_index in 0..2 {
let mut reader = PutObjReader::from_vec(format!("pool-{pool_index}-history").into_bytes());
store.pools[pool_index]
.put_object(
&bucket,
object,
&mut reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("historical version should be stored");
}
let marker = store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("trigger marker should be stored in the first pool");
let marker_id = marker.version_id.expect("trigger marker should have a version id");
let mut opts = ObjectOptions {
delete_prefix: true,
delete_prefix_object: true,
versioned: true,
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(marker_id),
delete_marker: true,
action: rustfs_common::metrics::IlmAction::DelMarkerDeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(
crate::bucket::metadata_sys::ObjectLockConfigState::ConfirmedAbsent,
))),
delete_replication_config_snapshot: Some(Arc::new(
crate::bucket::replication::DeleteReplicationConfigSnapshot::default(),
)),
..Default::default()
};
opts.ensure_lifecycle_delete_all_journal();
let _failure_guard = LifecycleDeleteAllFailureGuard;
*LIFECYCLE_DELETE_ALL_TEST_FAILURE
.lock()
.expect("lifecycle delete-all failure hook should not poison") =
Some((crate::object_api::LifecycleDeleteAllPhase::History, 1));
let err = store
.delete_prefix(&bucket, object, &opts)
.await
.expect_err("a later pool history failure must stop before trigger deletion");
assert_eq!(err, StorageError::PreconditionFailed);
assert!(
opts.lifecycle_delete_all_journal()
.expect("delete-all journal should be initialized")
.lock()
.mutation_started()
);
let first_pool = store.pools[0].disk_set[0]
.load_file_info_versions_exact(&bucket, object)
.await
.expect("first pool metadata should load")
.expect("the trigger should remain");
let first_pool_ordinary: Vec<&FileInfo> = first_pool
.versions
.iter()
.filter(|version| !version.tier_free_version())
.collect();
assert_eq!(first_pool_ordinary.len(), 1);
assert_eq!(first_pool_ordinary[0].version_id, Some(marker_id));
assert!(first_pool_ordinary[0].deleted);
*LIFECYCLE_DELETE_ALL_TEST_FAILURE
.lock()
.expect("lifecycle delete-all failure hook should not poison") = None;
store
.delete_prefix(&bucket, object, &opts)
.await
.expect("retry should delete remaining history and its trigger owner");
for pool in &store.pools {
assert!(
pool.disk_set[0]
.load_file_info_versions_exact(&bucket, object)
.await
.expect("pool metadata should load after retry")
.is_none(),
"all ordinary versions should be removed after retry"
);
}
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial]
async fn lifecycle_delete_all_phase_failures_preserve_barriers_and_retry() {
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("lifecycle-delete-all-phases", &[4, 4]).await;
let bucket = format!("lifecycle-delete-all-phases-{}", Uuid::new_v4().simple());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created in both pools");
let _failure_guard = LifecycleDeleteAllFailureGuard;
for (object, phase, expected_counts, mutation_started) in [
("preflight-failure", crate::object_api::LifecycleDeleteAllPhase::Preflight, [2, 2], false),
(
"final-preflight-failure",
crate::object_api::LifecycleDeleteAllPhase::FinalPreflight,
[1, 1],
true,
),
("trigger-failure", crate::object_api::LifecycleDeleteAllPhase::Trigger, [0, 1], true),
] {
let opts = seed_multi_pool_delete_all(&store, &bucket, object).await;
*LIFECYCLE_DELETE_ALL_TEST_FAILURE
.lock()
.expect("lifecycle delete-all failure hook should not poison") = Some((phase, 1));
let err = store
.delete_prefix(&bucket, object, &opts)
.await
.expect_err("injected phase failure should stop the transaction");
assert_eq!(err, StorageError::PreconditionFailed);
assert_eq!(
opts.lifecycle_delete_all_journal()
.expect("delete-all journal should be initialized")
.lock()
.mutation_started(),
mutation_started
);
assert_eq!(ordinary_version_count(&store, 0, &bucket, object).await, expected_counts[0]);
assert_eq!(ordinary_version_count(&store, 1, &bucket, object).await, expected_counts[1]);
*LIFECYCLE_DELETE_ALL_TEST_FAILURE
.lock()
.expect("lifecycle delete-all failure hook should not poison") = None;
store
.delete_prefix(&bucket, object, &opts)
.await
.expect("retry should converge after the injected failure is removed");
assert_eq!(ordinary_version_count(&store, 0, &bucket, object).await, 0);
assert_eq!(ordinary_version_count(&store, 1, &bucket, object).await, 0);
}
shutdown.cancel();
}
fn assert_backend_layout_empty(info: &rustfs_madmin::BackendInfo) {
assert!(info.standard_sc_parities.is_empty());
assert!(info.standard_sc_data.is_empty());