mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 20:36:38 +00:00
fix(ecstore): route rebalance delete markers to store and count lifecycle expiries (#4540)
fix(ecstore): route rebalance delete markers to store and count expiries Rebalance migrated delete markers via SetDisks::delete_object (single set, no cross-pool routing), which silently rewrote the marker back onto the source set. The subsequent source cleanup then deleted the whole entry, losing the delete marker and reviving logically-deleted objects (#942, P0). Route delete-marker migration through ECStore::delete_object via a store- level closure that mirrors the existing transfer closure, so the marker lands on the cross-pool target honouring data_movement/src_pool_idx/ delete_marker; on failure it is not counted as moved, so the source entry is not cleaned up. The unused MigrationBackend::delete_object_for_migration footgun is removed. Rebalance source cleanup also ignored lifecycle-expired versions: the gate used rebalanced == total_versions and passed an empty allowed_missing, so any entry with an expired version could never be cleaned up, leaking the migrated versions in the source pool (#950). Mirror decommission: gate on rebalanced + expired == total_versions and feed expired version identities into the cleanup preflight allowed_missing, plus a source_retained warning for parity with decommission diagnostics. Adds regression tests for delete-marker store routing and the expiry-aware cleanup gate. Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -33,8 +33,9 @@ use crate::core::pools::ListCallback;
|
|||||||
use crate::data_movement;
|
use crate::data_movement;
|
||||||
use crate::data_movement::backpressure::{self, DataMovementOperation};
|
use crate::data_movement::backpressure::{self, DataMovementOperation};
|
||||||
use crate::error::{Error, Result};
|
use crate::error::{Error, Result};
|
||||||
use crate::object_api::GetObjectReader;
|
use crate::object_api::{GetObjectReader, ObjectOptions};
|
||||||
use crate::set_disk::SetDisks;
|
use crate::set_disk::SetDisks;
|
||||||
|
use crate::storage_api_contracts::object::ObjectOperations as _;
|
||||||
use crate::store::ECStore;
|
use crate::store::ECStore;
|
||||||
use rustfs_filemeta::MetaCacheEntry;
|
use rustfs_filemeta::MetaCacheEntry;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
@@ -109,6 +110,7 @@ impl ECStore {
|
|||||||
|
|
||||||
let mut rebalanced: usize = 0;
|
let mut rebalanced: usize = 0;
|
||||||
let mut expired: usize = 0;
|
let mut expired: usize = 0;
|
||||||
|
let mut cleanup_preflight_allowed_missing = Vec::new();
|
||||||
let mut stats_updates = Vec::with_capacity(fivs.versions.len());
|
let mut stats_updates = Vec::with_capacity(fivs.versions.len());
|
||||||
for version in fivs.versions.iter() {
|
for version in fivs.versions.iter() {
|
||||||
if crate::core::pools::should_skip_lifecycle_for_data_movement(
|
if crate::core::pools::should_skip_lifecycle_for_data_movement(
|
||||||
@@ -123,6 +125,10 @@ impl ECStore {
|
|||||||
.await?
|
.await?
|
||||||
{
|
{
|
||||||
expired += 1;
|
expired += 1;
|
||||||
|
// The lifecycle expiry above physically deleted this version from the source set.
|
||||||
|
// Record its identity so the source-cleanup preflight tolerates its absence,
|
||||||
|
// mirroring decommission; otherwise the entry can never be cleaned up.
|
||||||
|
cleanup_preflight_allowed_missing.push(data_movement::source_cleanup_version_identity(version));
|
||||||
debug!(
|
debug!(
|
||||||
event = EVENT_REBALANCE_ENTRY,
|
event = EVENT_REBALANCE_ENTRY,
|
||||||
component = LOG_COMPONENT_ECSTORE,
|
component = LOG_COMPONENT_ECSTORE,
|
||||||
@@ -159,6 +165,12 @@ impl ECStore {
|
|||||||
let store = self.clone();
|
let store = self.clone();
|
||||||
async move { store.rebalance_object(src_pool_idx, bucket, rd).await }
|
async move { store.rebalance_object(src_pool_idx, bucket, rd).await }
|
||||||
};
|
};
|
||||||
|
// Route delete-marker migration through the store layer so it lands on the
|
||||||
|
// cross-pool target (excluding the source pool), not back onto the source set.
|
||||||
|
let mut delete_marker = |bucket: String, object: String, opts: ObjectOptions| {
|
||||||
|
let store = self.clone();
|
||||||
|
async move { store.delete_object(&bucket, &object, opts).await }
|
||||||
|
};
|
||||||
let result = migrate_entry_version(
|
let result = migrate_entry_version(
|
||||||
set.as_ref(),
|
set.as_ref(),
|
||||||
bucket.clone(),
|
bucket.clone(),
|
||||||
@@ -168,6 +180,7 @@ impl ECStore {
|
|||||||
rebalance_max_attempts(),
|
rebalance_max_attempts(),
|
||||||
should_ignore_rebalance_data_usage_cache(bucket.as_str()),
|
should_ignore_rebalance_data_usage_cache(bucket.as_str()),
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
&mut delete_marker,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -255,14 +268,14 @@ impl ECStore {
|
|||||||
entry.name.as_str(),
|
entry.name.as_str(),
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
if should_cleanup_rebalance_source_entry(rebalanced, fivs.versions.len()) {
|
if should_cleanup_rebalance_source_entry(rebalanced, fivs.versions.len(), expired) {
|
||||||
let cleanup_warning = resolve_rebalance_entry_cleanup_delete_result(
|
let cleanup_warning = resolve_rebalance_entry_cleanup_delete_result(
|
||||||
data_movement::cleanup_source_entry_if_unchanged(
|
data_movement::cleanup_source_entry_if_unchanged(
|
||||||
set.clone(),
|
set.clone(),
|
||||||
bucket.as_str(),
|
bucket.as_str(),
|
||||||
entry.name.as_str(),
|
entry.name.as_str(),
|
||||||
&fivs,
|
&fivs,
|
||||||
&[],
|
&cleanup_preflight_allowed_missing,
|
||||||
"rebalance",
|
"rebalance",
|
||||||
)
|
)
|
||||||
.await,
|
.await,
|
||||||
@@ -310,6 +323,20 @@ impl ECStore {
|
|||||||
"Deleted rebalance source entry"
|
"Deleted rebalance source entry"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
} else if rebalanced != fivs.versions.len() || expired > 0 {
|
||||||
|
warn!(
|
||||||
|
event = EVENT_REBALANCE_ENTRY,
|
||||||
|
component = LOG_COMPONENT_ECSTORE,
|
||||||
|
subsystem = LOG_SUBSYSTEM_REBALANCE,
|
||||||
|
pool_index,
|
||||||
|
bucket = %bucket,
|
||||||
|
object = %entry.name,
|
||||||
|
rebalanced,
|
||||||
|
total_versions = fivs.versions.len(),
|
||||||
|
expired,
|
||||||
|
state = "source_retained",
|
||||||
|
"Rebalance source object retained"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(RebalanceEntryOutcome::Completed)
|
Ok(RebalanceEntryOutcome::Completed)
|
||||||
|
|||||||
@@ -4,10 +4,7 @@ use crate::data_usage::DATA_USAGE_CACHE_NAME;
|
|||||||
use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found};
|
use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found};
|
||||||
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions};
|
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions};
|
||||||
use crate::set_disk::SetDisks;
|
use crate::set_disk::SetDisks;
|
||||||
use crate::storage_api_contracts::{
|
use crate::storage_api_contracts::{object::ObjectIO, range::HTTPRangeSpec};
|
||||||
object::{ObjectIO, ObjectOperations as _},
|
|
||||||
range::HTTPRangeSpec,
|
|
||||||
};
|
|
||||||
use http::HeaderMap;
|
use http::HeaderMap;
|
||||||
use rustfs_filemeta::FileInfo;
|
use rustfs_filemeta::FileInfo;
|
||||||
use rustfs_utils::path::encode_dir_object;
|
use rustfs_utils::path::encode_dir_object;
|
||||||
@@ -64,8 +61,6 @@ pub(crate) trait MigrationBackend: Send + Sync {
|
|||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
) -> Result<GetObjectReader>;
|
) -> Result<GetObjectReader>;
|
||||||
|
|
||||||
async fn delete_object_for_migration(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result<ObjectInfo>;
|
|
||||||
|
|
||||||
async fn move_remote_version_for_migration(
|
async fn move_remote_version_for_migration(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -88,10 +83,6 @@ impl MigrationBackend for SetDisks {
|
|||||||
self.get_object_reader(bucket, object, range, h, opts).await
|
self.get_object_reader(bucket, object, range, h, opts).await
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_object_for_migration(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result<ObjectInfo> {
|
|
||||||
self.delete_object(bucket, object, opts).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn move_remote_version_for_migration(
|
async fn move_remote_version_for_migration(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -104,7 +95,7 @@ impl MigrationBackend for SetDisks {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[allow(clippy::too_many_arguments)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
pub(crate) async fn migrate_entry_version<Backend, F, Fut>(
|
pub(crate) async fn migrate_entry_version<Backend, F, Fut, D, DFut>(
|
||||||
set: &Backend,
|
set: &Backend,
|
||||||
bucket: String,
|
bucket: String,
|
||||||
pool_index: usize,
|
pool_index: usize,
|
||||||
@@ -113,11 +104,14 @@ pub(crate) async fn migrate_entry_version<Backend, F, Fut>(
|
|||||||
max_attempts: usize,
|
max_attempts: usize,
|
||||||
ignore_data_usage_cache: bool,
|
ignore_data_usage_cache: bool,
|
||||||
transfer: F,
|
transfer: F,
|
||||||
|
delete_marker: D,
|
||||||
) -> MigrationVersionResult
|
) -> MigrationVersionResult
|
||||||
where
|
where
|
||||||
Backend: MigrationBackend + ?Sized,
|
Backend: MigrationBackend + ?Sized,
|
||||||
F: FnMut(usize, String, GetObjectReader) -> Fut + Send,
|
F: FnMut(usize, String, GetObjectReader) -> Fut + Send,
|
||||||
Fut: Future<Output = Result<()>> + Send,
|
Fut: Future<Output = Result<()>> + Send,
|
||||||
|
D: FnMut(String, String, ObjectOptions) -> DFut + Send,
|
||||||
|
DFut: Future<Output = Result<ObjectInfo>> + Send,
|
||||||
{
|
{
|
||||||
migrate_entry_version_with_retry_wait(
|
migrate_entry_version_with_retry_wait(
|
||||||
set,
|
set,
|
||||||
@@ -128,13 +122,14 @@ where
|
|||||||
max_attempts,
|
max_attempts,
|
||||||
ignore_data_usage_cache,
|
ignore_data_usage_cache,
|
||||||
transfer,
|
transfer,
|
||||||
|
delete_marker,
|
||||||
sleep_rebalance_migration_retry,
|
sleep_rebalance_migration_retry,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(clippy::too_many_arguments)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
pub(super) async fn migrate_entry_version_with_retry_wait<Backend, F, Fut, W, WFut>(
|
pub(super) async fn migrate_entry_version_with_retry_wait<Backend, F, Fut, D, DFut, W, WFut>(
|
||||||
set: &Backend,
|
set: &Backend,
|
||||||
bucket: String,
|
bucket: String,
|
||||||
pool_index: usize,
|
pool_index: usize,
|
||||||
@@ -143,12 +138,15 @@ pub(super) async fn migrate_entry_version_with_retry_wait<Backend, F, Fut, W, WF
|
|||||||
max_attempts: usize,
|
max_attempts: usize,
|
||||||
ignore_data_usage_cache: bool,
|
ignore_data_usage_cache: bool,
|
||||||
mut transfer: F,
|
mut transfer: F,
|
||||||
|
mut delete_marker: D,
|
||||||
mut wait_retry: W,
|
mut wait_retry: W,
|
||||||
) -> MigrationVersionResult
|
) -> MigrationVersionResult
|
||||||
where
|
where
|
||||||
Backend: MigrationBackend + ?Sized,
|
Backend: MigrationBackend + ?Sized,
|
||||||
F: FnMut(usize, String, GetObjectReader) -> Fut + Send,
|
F: FnMut(usize, String, GetObjectReader) -> Fut + Send,
|
||||||
Fut: Future<Output = Result<()>> + Send,
|
Fut: Future<Output = Result<()>> + Send,
|
||||||
|
D: FnMut(String, String, ObjectOptions) -> DFut + Send,
|
||||||
|
DFut: Future<Output = Result<ObjectInfo>> + Send,
|
||||||
W: FnMut(Duration) -> WFut + Send,
|
W: FnMut(Duration) -> WFut + Send,
|
||||||
WFut: Future<Output = ()> + Send,
|
WFut: Future<Output = ()> + Send,
|
||||||
{
|
{
|
||||||
@@ -207,9 +205,16 @@ where
|
|||||||
}
|
}
|
||||||
|
|
||||||
if version.deleted {
|
if version.deleted {
|
||||||
if let Err(err) = set
|
// Delete markers must be routed through the store layer (ECStore::delete_object /
|
||||||
.delete_object_for_migration(&bucket, &version.name, rebalance_delete_marker_opts(version, version_id, pool_index))
|
// handle_delete_object), which honours data_movement/src_pool_idx/delete_marker and
|
||||||
.await
|
// writes the marker to the cross-pool target. Writing via the source SetDisks would
|
||||||
|
// silently rewrite the marker back onto the source set and lose it during cleanup.
|
||||||
|
if let Err(err) = delete_marker(
|
||||||
|
bucket.clone(),
|
||||||
|
version.name.clone(),
|
||||||
|
rebalance_delete_marker_opts(version, version_id, pool_index),
|
||||||
|
)
|
||||||
|
.await
|
||||||
{
|
{
|
||||||
if is_err_object_not_found(&err) || is_err_version_not_found(&err) {
|
if is_err_object_not_found(&err) || is_err_version_not_found(&err) {
|
||||||
return MigrationVersionResult {
|
return MigrationVersionResult {
|
||||||
|
|||||||
@@ -110,25 +110,20 @@ struct LegacyRebalanceMeta {
|
|||||||
|
|
||||||
struct MigrationBackendSpy {
|
struct MigrationBackendSpy {
|
||||||
get_object_reader: Mutex<Option<core::result::Result<GetObjectReader, Error>>>,
|
get_object_reader: Mutex<Option<core::result::Result<GetObjectReader, Error>>>,
|
||||||
delete_object: Mutex<Option<core::result::Result<ObjectInfo, Error>>>,
|
|
||||||
move_remote: Mutex<Option<core::result::Result<(), Error>>>,
|
move_remote: Mutex<Option<core::result::Result<(), Error>>>,
|
||||||
get_calls: AtomicUsize,
|
get_calls: AtomicUsize,
|
||||||
delete_calls: AtomicUsize,
|
|
||||||
move_remote_calls: AtomicUsize,
|
move_remote_calls: AtomicUsize,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl MigrationBackendSpy {
|
impl MigrationBackendSpy {
|
||||||
fn new(
|
fn new(
|
||||||
get_object_reader: Option<core::result::Result<GetObjectReader, Error>>,
|
get_object_reader: Option<core::result::Result<GetObjectReader, Error>>,
|
||||||
delete_object: Option<core::result::Result<ObjectInfo, Error>>,
|
|
||||||
move_remote: Option<core::result::Result<(), Error>>,
|
move_remote: Option<core::result::Result<(), Error>>,
|
||||||
) -> Self {
|
) -> Self {
|
||||||
Self {
|
Self {
|
||||||
get_object_reader: Mutex::new(get_object_reader),
|
get_object_reader: Mutex::new(get_object_reader),
|
||||||
delete_object: Mutex::new(delete_object),
|
|
||||||
move_remote: Mutex::new(move_remote),
|
move_remote: Mutex::new(move_remote),
|
||||||
get_calls: AtomicUsize::new(0),
|
get_calls: AtomicUsize::new(0),
|
||||||
delete_calls: AtomicUsize::new(0),
|
|
||||||
move_remote_calls: AtomicUsize::new(0),
|
move_remote_calls: AtomicUsize::new(0),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -137,10 +132,6 @@ impl MigrationBackendSpy {
|
|||||||
self.get_calls.load(Ordering::SeqCst)
|
self.get_calls.load(Ordering::SeqCst)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn delete_calls(&self) -> usize {
|
|
||||||
self.delete_calls.load(Ordering::SeqCst)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn move_remote_calls(&self) -> usize {
|
fn move_remote_calls(&self) -> usize {
|
||||||
self.move_remote_calls.load(Ordering::SeqCst)
|
self.move_remote_calls.load(Ordering::SeqCst)
|
||||||
}
|
}
|
||||||
@@ -172,15 +163,6 @@ impl MigrationBackend for MigrationBackendSpy {
|
|||||||
Ok(Self::make_reader())
|
Ok(Self::make_reader())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_object_for_migration(&self, _bucket: &str, _object: &str, _opts: ObjectOptions) -> Result<ObjectInfo> {
|
|
||||||
self.delete_calls.fetch_add(1, Ordering::SeqCst);
|
|
||||||
if let Some(result) = self.delete_object.lock().unwrap().take() {
|
|
||||||
return result;
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(ObjectInfo::default())
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn move_remote_version_for_migration(
|
async fn move_remote_version_for_migration(
|
||||||
&self,
|
&self,
|
||||||
_bucket: &str,
|
_bucket: &str,
|
||||||
@@ -249,7 +231,7 @@ fn test_rebalance_delete_marker_opts_preserves_replication_state() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
|
async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
|
||||||
let backend = MigrationBackendSpy::new(None, Some(Ok(ObjectInfo::default())), Some(Ok(())));
|
let backend = MigrationBackendSpy::new(None, Some(Ok(())));
|
||||||
let version = version_remote();
|
let version = version_remote();
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
@@ -272,6 +254,7 @@ async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -283,16 +266,12 @@ async fn test_migrate_entry_version_remote_version_is_moved_without_transfer() {
|
|||||||
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
||||||
assert_eq!(backend.move_remote_calls(), 1);
|
assert_eq!(backend.move_remote_calls(), 1);
|
||||||
assert_eq!(backend.get_calls(), 0);
|
assert_eq!(backend.get_calls(), 0);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_remote_not_found_is_cleanup_ignored() {
|
async fn test_migrate_entry_version_remote_not_found_is_cleanup_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(
|
let backend =
|
||||||
None,
|
MigrationBackendSpy::new(None, Some(Err(Error::ObjectNotFound("bucket".to_string(), "object.bin".to_string()))));
|
||||||
Some(Ok(ObjectInfo::default())),
|
|
||||||
Some(Err(Error::ObjectNotFound("bucket".to_string(), "object.bin".to_string()))),
|
|
||||||
);
|
|
||||||
let version = version_remote();
|
let version = version_remote();
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
@@ -315,6 +294,7 @@ async fn test_migrate_entry_version_remote_not_found_is_cleanup_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -326,13 +306,11 @@ async fn test_migrate_entry_version_remote_not_found_is_cleanup_ignored() {
|
|||||||
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
||||||
assert_eq!(backend.move_remote_calls(), 1);
|
assert_eq!(backend.move_remote_calls(), 1);
|
||||||
assert_eq!(backend.get_calls(), 0);
|
assert_eq!(backend.get_calls(), 0);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_remote_overwrite_is_not_ignored() {
|
async fn test_migrate_entry_version_remote_overwrite_is_not_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(
|
let backend = MigrationBackendSpy::new(
|
||||||
None,
|
|
||||||
None,
|
None,
|
||||||
Some(Err(Error::DataMovementOverwriteErr(
|
Some(Err(Error::DataMovementOverwriteErr(
|
||||||
"bucket".to_string(),
|
"bucket".to_string(),
|
||||||
@@ -352,6 +330,7 @@ async fn test_migrate_entry_version_remote_overwrite_is_not_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -366,7 +345,7 @@ async fn test_migrate_entry_version_remote_overwrite_is_not_ignored() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_remote_failure_is_reported() {
|
async fn test_migrate_entry_version_remote_failure_is_reported() {
|
||||||
let backend = MigrationBackendSpy::new(None, Some(Ok(ObjectInfo::default())), Some(Err(Error::SlowDown)));
|
let backend = MigrationBackendSpy::new(None, Some(Err(Error::SlowDown)));
|
||||||
let version = version_remote();
|
let version = version_remote();
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
@@ -389,6 +368,7 @@ async fn test_migrate_entry_version_remote_failure_is_reported() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -400,14 +380,26 @@ async fn test_migrate_entry_version_remote_failure_is_reported() {
|
|||||||
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
||||||
assert_eq!(backend.move_remote_calls(), 1);
|
assert_eq!(backend.move_remote_calls(), 1);
|
||||||
assert_eq!(backend.get_calls(), 0);
|
assert_eq!(backend.get_calls(), 0);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_deleted_version_calls_delete_and_moved() {
|
async fn test_migrate_entry_version_deleted_version_routes_delete_through_store_and_moved() {
|
||||||
let backend = MigrationBackendSpy::new(None, Some(Ok(ObjectInfo::default())), None);
|
let backend = MigrationBackendSpy::new(None, None);
|
||||||
let version = version_deleted();
|
let version = version_deleted();
|
||||||
let mut transfer = |_, _, _| async move { Ok(()) };
|
let mut transfer = |_, _, _| async move { Ok(()) };
|
||||||
|
// The delete marker must be routed through the store closure (cross-pool routing), never
|
||||||
|
// through the source SetDisks. Assert the closure is invoked and the source set is not.
|
||||||
|
let delete_calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let mut delete_marker = {
|
||||||
|
let delete_calls = delete_calls.clone();
|
||||||
|
move |_: String, _: String, _: ObjectOptions| {
|
||||||
|
let delete_calls = delete_calls.clone();
|
||||||
|
async move {
|
||||||
|
delete_calls.fetch_add(1, Ordering::SeqCst);
|
||||||
|
Ok(ObjectInfo::default())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let result = migrate_entry_version(
|
let result = migrate_entry_version(
|
||||||
&backend,
|
&backend,
|
||||||
@@ -418,6 +410,7 @@ async fn test_migrate_entry_version_deleted_version_calls_delete_and_moved() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
&mut delete_marker,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -427,18 +420,25 @@ async fn test_migrate_entry_version_deleted_version_calls_delete_and_moved() {
|
|||||||
assert!(!result.failed);
|
assert!(!result.failed);
|
||||||
assert!(result.error.is_none());
|
assert!(result.error.is_none());
|
||||||
assert_eq!(backend.get_calls(), 0);
|
assert_eq!(backend.get_calls(), 0);
|
||||||
assert_eq!(backend.delete_calls(), 1);
|
assert_eq!(delete_calls.load(Ordering::SeqCst), 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_deleted_version_not_found_is_ignored() {
|
async fn test_migrate_entry_version_deleted_version_not_found_is_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(
|
let backend = MigrationBackendSpy::new(None, None);
|
||||||
None,
|
|
||||||
Some(Err(Error::ObjectNotFound("bucket".to_string(), "object.bin".to_string()))),
|
|
||||||
None,
|
|
||||||
);
|
|
||||||
let version = version_deleted();
|
let version = version_deleted();
|
||||||
let mut transfer = |_, _, _| async move { Ok(()) };
|
let mut transfer = |_, _, _| async move { Ok(()) };
|
||||||
|
let delete_calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let mut delete_marker = {
|
||||||
|
let delete_calls = delete_calls.clone();
|
||||||
|
move |_: String, _: String, _: ObjectOptions| {
|
||||||
|
let delete_calls = delete_calls.clone();
|
||||||
|
async move {
|
||||||
|
delete_calls.fetch_add(1, Ordering::SeqCst);
|
||||||
|
Err(Error::ObjectNotFound("bucket".to_string(), "object.bin".to_string()))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let result = migrate_entry_version(
|
let result = migrate_entry_version(
|
||||||
&backend,
|
&backend,
|
||||||
@@ -449,6 +449,7 @@ async fn test_migrate_entry_version_deleted_version_not_found_is_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
&mut delete_marker,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -457,22 +458,29 @@ async fn test_migrate_entry_version_deleted_version_not_found_is_ignored() {
|
|||||||
assert!(!result.moved);
|
assert!(!result.moved);
|
||||||
assert!(!result.failed);
|
assert!(!result.failed);
|
||||||
assert!(result.error.is_none());
|
assert!(result.error.is_none());
|
||||||
assert_eq!(backend.delete_calls(), 1);
|
assert_eq!(delete_calls.load(Ordering::SeqCst), 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_deleted_version_overwrite_is_not_ignored() {
|
async fn test_migrate_entry_version_deleted_version_overwrite_is_not_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(
|
let backend = MigrationBackendSpy::new(None, None);
|
||||||
None,
|
|
||||||
Some(Err(Error::DataMovementOverwriteErr(
|
|
||||||
"bucket".to_string(),
|
|
||||||
"object.bin".to_string(),
|
|
||||||
"vid-1".to_string(),
|
|
||||||
))),
|
|
||||||
None,
|
|
||||||
);
|
|
||||||
let version = version_deleted();
|
let version = version_deleted();
|
||||||
let mut transfer = |_, _, _| async move { Ok(()) };
|
let mut transfer = |_, _, _| async move { Ok(()) };
|
||||||
|
let delete_calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let mut delete_marker = {
|
||||||
|
let delete_calls = delete_calls.clone();
|
||||||
|
move |_: String, _: String, _: ObjectOptions| {
|
||||||
|
let delete_calls = delete_calls.clone();
|
||||||
|
async move {
|
||||||
|
delete_calls.fetch_add(1, Ordering::SeqCst);
|
||||||
|
Err(Error::DataMovementOverwriteErr(
|
||||||
|
"bucket".to_string(),
|
||||||
|
"object.bin".to_string(),
|
||||||
|
"vid-1".to_string(),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let result = migrate_entry_version(
|
let result = migrate_entry_version(
|
||||||
&backend,
|
&backend,
|
||||||
@@ -483,6 +491,7 @@ async fn test_migrate_entry_version_deleted_version_overwrite_is_not_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
&mut delete_marker,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -492,16 +501,13 @@ async fn test_migrate_entry_version_deleted_version_overwrite_is_not_ignored() {
|
|||||||
assert!(!result.moved);
|
assert!(!result.moved);
|
||||||
assert_eq!(result.stage, Some("delete_marker"));
|
assert_eq!(result.stage, Some("delete_marker"));
|
||||||
assert!(matches!(result.error, Some(Error::DataMovementOverwriteErr(_, _, _))));
|
assert!(matches!(result.error, Some(Error::DataMovementOverwriteErr(_, _, _))));
|
||||||
assert_eq!(backend.delete_calls(), 1);
|
assert_eq!(delete_calls.load(Ordering::SeqCst), 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_reader_not_found_is_ignored() {
|
async fn test_migrate_entry_version_reader_not_found_is_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(
|
let backend =
|
||||||
Some(Err(Error::ObjectNotFound("bucket".to_string(), "object.bin".to_string()))),
|
MigrationBackendSpy::new(Some(Err(Error::ObjectNotFound("bucket".to_string(), "object.bin".to_string()))), None);
|
||||||
None,
|
|
||||||
None,
|
|
||||||
);
|
|
||||||
let version = version_normal();
|
let version = version_normal();
|
||||||
let mut transfer = |_, _, _| async move { Ok(()) };
|
let mut transfer = |_, _, _| async move { Ok(()) };
|
||||||
|
|
||||||
@@ -514,6 +520,7 @@ async fn test_migrate_entry_version_reader_not_found_is_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -523,12 +530,11 @@ async fn test_migrate_entry_version_reader_not_found_is_ignored() {
|
|||||||
assert!(!result.failed);
|
assert!(!result.failed);
|
||||||
assert!(result.error.is_none());
|
assert!(result.error.is_none());
|
||||||
assert_eq!(backend.get_calls(), 1);
|
assert_eq!(backend.get_calls(), 1);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_reader_retries_before_success() {
|
async fn test_migrate_entry_version_reader_retries_before_success() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Err(Error::SlowDown)), None, None);
|
let backend = MigrationBackendSpy::new(Some(Err(Error::SlowDown)), None);
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let wait_count = Arc::new(AtomicUsize::new(0));
|
let wait_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
@@ -552,6 +558,7 @@ async fn test_migrate_entry_version_reader_retries_before_success() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
{
|
{
|
||||||
let wait_count = wait_count.clone();
|
let wait_count = wait_count.clone();
|
||||||
move |_| {
|
move |_| {
|
||||||
@@ -570,7 +577,6 @@ async fn test_migrate_entry_version_reader_retries_before_success() {
|
|||||||
assert!(!result.failed);
|
assert!(!result.failed);
|
||||||
assert!(result.error.is_none());
|
assert!(result.error.is_none());
|
||||||
assert_eq!(backend.get_calls(), 2);
|
assert_eq!(backend.get_calls(), 2);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
assert_eq!(transfer_count.load(Ordering::SeqCst), 1);
|
assert_eq!(transfer_count.load(Ordering::SeqCst), 1);
|
||||||
assert_eq!(wait_count.load(Ordering::SeqCst), 1);
|
assert_eq!(wait_count.load(Ordering::SeqCst), 1);
|
||||||
}
|
}
|
||||||
@@ -605,10 +611,6 @@ impl MigrationBackend for AlwaysFailGetBackend {
|
|||||||
Err(Error::SlowDown)
|
Err(Error::SlowDown)
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_object_for_migration(&self, _bucket: &str, _object: &str, _opts: ObjectOptions) -> Result<ObjectInfo> {
|
|
||||||
Ok(ObjectInfo::default())
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn move_remote_version_for_migration(
|
async fn move_remote_version_for_migration(
|
||||||
&self,
|
&self,
|
||||||
_bucket: &str,
|
_bucket: &str,
|
||||||
@@ -645,6 +647,7 @@ async fn test_migrate_entry_version_reader_fails_after_retries() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -682,6 +685,7 @@ async fn test_migrate_entry_version_zero_max_attempts_still_attempts_once() {
|
|||||||
0,
|
0,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -696,7 +700,7 @@ async fn test_migrate_entry_version_zero_max_attempts_still_attempts_once() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_transfer_retries_before_success() {
|
async fn test_migrate_entry_version_transfer_retries_before_success() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let wait_count = Arc::new(AtomicUsize::new(0));
|
let wait_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
@@ -723,6 +727,7 @@ async fn test_migrate_entry_version_transfer_retries_before_success() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
{
|
{
|
||||||
let wait_count = wait_count.clone();
|
let wait_count = wait_count.clone();
|
||||||
move |_| {
|
move |_| {
|
||||||
@@ -746,7 +751,7 @@ async fn test_migrate_entry_version_transfer_retries_before_success() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_transfer_non_transient_fails_without_retry() {
|
async fn test_migrate_entry_version_transfer_non_transient_fails_without_retry() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let wait_count = Arc::new(AtomicUsize::new(0));
|
let wait_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
@@ -770,6 +775,7 @@ async fn test_migrate_entry_version_transfer_non_transient_fails_without_retry()
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
{
|
{
|
||||||
let wait_count = wait_count.clone();
|
let wait_count = wait_count.clone();
|
||||||
move |_| {
|
move |_| {
|
||||||
@@ -793,7 +799,7 @@ async fn test_migrate_entry_version_transfer_non_transient_fails_without_retry()
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_transfer_fails_after_retries() {
|
async fn test_migrate_entry_version_transfer_fails_after_retries() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
let transfer_count = transfer_count.clone();
|
let transfer_count = transfer_count.clone();
|
||||||
@@ -816,6 +822,7 @@ async fn test_migrate_entry_version_transfer_fails_after_retries() {
|
|||||||
2,
|
2,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -830,7 +837,7 @@ async fn test_migrate_entry_version_transfer_fails_after_retries() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_transfer_not_found_is_ignored() {
|
async fn test_migrate_entry_version_transfer_not_found_is_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
let transfer_count = transfer_count.clone();
|
let transfer_count = transfer_count.clone();
|
||||||
@@ -853,6 +860,7 @@ async fn test_migrate_entry_version_transfer_not_found_is_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -866,7 +874,7 @@ async fn test_migrate_entry_version_transfer_not_found_is_ignored() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_transfer_overwrite_is_not_ignored() {
|
async fn test_migrate_entry_version_transfer_overwrite_is_not_ignored() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let transfer_count = Arc::new(AtomicUsize::new(0));
|
let transfer_count = Arc::new(AtomicUsize::new(0));
|
||||||
let mut transfer = {
|
let mut transfer = {
|
||||||
let transfer_count = transfer_count.clone();
|
let transfer_count = transfer_count.clone();
|
||||||
@@ -893,6 +901,7 @@ async fn test_migrate_entry_version_transfer_overwrite_is_not_ignored() {
|
|||||||
3,
|
3,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -907,7 +916,7 @@ async fn test_migrate_entry_version_transfer_overwrite_is_not_ignored() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_ignores_data_usage_cache_when_enabled() {
|
async fn test_migrate_entry_version_ignores_data_usage_cache_when_enabled() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let version = {
|
let version = {
|
||||||
let mut version = version_normal();
|
let mut version = version_normal();
|
||||||
version.name = format!("{}.{}", DATA_USAGE_CACHE_NAME, version.name);
|
version.name = format!("{}.{}", DATA_USAGE_CACHE_NAME, version.name);
|
||||||
@@ -934,6 +943,7 @@ async fn test_migrate_entry_version_ignores_data_usage_cache_when_enabled() {
|
|||||||
2,
|
2,
|
||||||
true,
|
true,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -944,12 +954,11 @@ async fn test_migrate_entry_version_ignores_data_usage_cache_when_enabled() {
|
|||||||
assert!(result.error.is_none());
|
assert!(result.error.is_none());
|
||||||
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
assert_eq!(transfer_count.load(Ordering::SeqCst), 0);
|
||||||
assert_eq!(backend.get_calls(), 0);
|
assert_eq!(backend.get_calls(), 0);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_data_usage_cache_moves_when_ignore_disabled() {
|
async fn test_migrate_entry_version_data_usage_cache_moves_when_ignore_disabled() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let version = {
|
let version = {
|
||||||
let mut version = version_normal();
|
let mut version = version_normal();
|
||||||
version.name = format!("{}.{}", DATA_USAGE_CACHE_NAME, version.name);
|
version.name = format!("{}.{}", DATA_USAGE_CACHE_NAME, version.name);
|
||||||
@@ -976,6 +985,7 @@ async fn test_migrate_entry_version_data_usage_cache_moves_when_ignore_disabled(
|
|||||||
2,
|
2,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -986,7 +996,6 @@ async fn test_migrate_entry_version_data_usage_cache_moves_when_ignore_disabled(
|
|||||||
assert!(result.error.is_none());
|
assert!(result.error.is_none());
|
||||||
assert_eq!(transfer_count.load(Ordering::SeqCst), 1);
|
assert_eq!(transfer_count.load(Ordering::SeqCst), 1);
|
||||||
assert_eq!(backend.get_calls(), 1);
|
assert_eq!(backend.get_calls(), 1);
|
||||||
assert_eq!(backend.delete_calls(), 0);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -1864,7 +1873,7 @@ fn test_with_rebalance_entry_context_formats_precise_stage() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_migrate_entry_version_transfer_failure_reports_write_target_stage() {
|
async fn test_migrate_entry_version_transfer_failure_reports_write_target_stage() {
|
||||||
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None, None);
|
let backend = MigrationBackendSpy::new(Some(Ok(MigrationBackendSpy::make_reader())), None);
|
||||||
let mut transfer = |_, _, _| async { Err(Error::SlowDown) };
|
let mut transfer = |_, _, _| async { Err(Error::SlowDown) };
|
||||||
let version = version_normal();
|
let version = version_normal();
|
||||||
|
|
||||||
@@ -1877,6 +1886,7 @@ async fn test_migrate_entry_version_transfer_failure_reports_write_target_stage(
|
|||||||
1,
|
1,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -1900,6 +1910,7 @@ async fn test_migrate_entry_version_reader_failure_reports_read_source_stage() {
|
|||||||
1,
|
1,
|
||||||
false,
|
false,
|
||||||
&mut transfer,
|
&mut transfer,
|
||||||
|
|_: String, _: String, _: ObjectOptions| async move { Ok::<_, Error>(ObjectInfo::default()) },
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -2026,12 +2037,35 @@ fn test_should_skip_rebalance_delete_marker_rejects_multiple_remaining_versions(
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_should_cleanup_rebalance_source_entry_accepts_all_versions_completed() {
|
fn test_should_cleanup_rebalance_source_entry_accepts_all_versions_completed() {
|
||||||
assert!(should_cleanup_rebalance_source_entry(3, 3));
|
assert!(should_cleanup_rebalance_source_entry(3, 3, 0));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_should_cleanup_rebalance_source_entry_rejects_versions_only_expired_by_lifecycle() {
|
fn test_should_cleanup_rebalance_source_entry_accepts_migrated_and_safely_expired_versions() {
|
||||||
assert!(!should_cleanup_rebalance_source_entry(2, 3));
|
// A multi-version object where one version was migrated and another expired by lifecycle
|
||||||
|
// must still be cleaned up; otherwise the migrated version leaks in the source pool.
|
||||||
|
assert!(should_cleanup_rebalance_source_entry(1, 2, 1));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_should_cleanup_rebalance_source_entry_accepts_single_version_only_expired_by_lifecycle() {
|
||||||
|
assert!(should_cleanup_rebalance_source_entry(0, 1, 1));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_should_cleanup_rebalance_source_entry_accepts_versions_only_expired_by_lifecycle() {
|
||||||
|
assert!(should_cleanup_rebalance_source_entry(0, 2, 2));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_should_cleanup_rebalance_source_entry_rejects_unmigrated_version() {
|
||||||
|
// One version neither migrated nor expired must block source cleanup.
|
||||||
|
assert!(!should_cleanup_rebalance_source_entry(1, 2, 0));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_should_cleanup_rebalance_source_entry_rejects_counter_overrun() {
|
||||||
|
assert!(!should_cleanup_rebalance_source_entry(2, 2, 1));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
@@ -347,8 +347,8 @@ pub(super) fn should_count_rebalance_version_complete(result: &MigrationVersionR
|
|||||||
result.cleanup_ignored || (result.moved && !result.failed)
|
result.cleanup_ignored || (result.moved && !result.failed)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(super) fn should_cleanup_rebalance_source_entry(rebalanced: usize, total_versions: usize) -> bool {
|
pub(super) fn should_cleanup_rebalance_source_entry(rebalanced: usize, total_versions: usize, expired: usize) -> bool {
|
||||||
rebalanced == total_versions
|
rebalanced.saturating_add(expired) == total_versions
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(super) fn should_skip_rebalance_delete_marker(
|
pub(super) fn should_skip_rebalance_delete_marker(
|
||||||
|
|||||||
Reference in New Issue
Block a user