mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 04:16:38 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2c3e68ad89 | |||
| 2f0918f60b | |||
| 5b951de2b7 |
@@ -400,7 +400,7 @@ jobs:
|
|||||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||||
needs: [ quick-checks ]
|
needs: [ quick-checks ]
|
||||||
runs-on: sm-standard-4
|
runs-on: sm-standard-4
|
||||||
timeout-minutes: 45
|
timeout-minutes: 90
|
||||||
env:
|
env:
|
||||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||||
steps:
|
steps:
|
||||||
@@ -440,7 +440,7 @@ jobs:
|
|||||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||||
needs: [ quick-checks ]
|
needs: [ quick-checks ]
|
||||||
runs-on: sm-standard-4
|
runs-on: sm-standard-4
|
||||||
timeout-minutes: 60
|
timeout-minutes: 90
|
||||||
env:
|
env:
|
||||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||||
steps:
|
steps:
|
||||||
@@ -470,7 +470,7 @@ jobs:
|
|||||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||||
needs: [ quick-checks ]
|
needs: [ quick-checks ]
|
||||||
runs-on: sm-standard-4
|
runs-on: sm-standard-4
|
||||||
timeout-minutes: 60
|
timeout-minutes: 90
|
||||||
strategy:
|
strategy:
|
||||||
# On a PR, one failing protocol leg is enough to know the PR is not ready,
|
# On a PR, one failing protocol leg is enough to know the PR is not ready,
|
||||||
# so stop the sibling leg instead of paying another ~40 minutes for it.
|
# so stop the sibling leg instead of paying another ~40 minutes for it.
|
||||||
|
|||||||
@@ -57,6 +57,13 @@ pub const DEFAULT_MAX_IO_EVENTS_PER_TICK: usize = 1024;
|
|||||||
pub const DEFAULT_EVENT_INTERVAL: u32 = 61;
|
pub const DEFAULT_EVENT_INTERVAL: u32 = 61;
|
||||||
pub const DEFAULT_RNG_SEED: Option<u64> = None; // None means random
|
pub const DEFAULT_RNG_SEED: Option<u64> = None; // None means random
|
||||||
|
|
||||||
|
/// Dedicated blocking thread pool for fsync/fdatasync operations.
|
||||||
|
/// When > 1, fsync operations are isolated from the main blocking pool to
|
||||||
|
/// prevent device-bound fsync from starving read operations (pread/stat/open).
|
||||||
|
/// Default 0 means auto (no isolation, use main runtime).
|
||||||
|
pub const ENV_FSYNC_BLOCKING_THREADS: &str = "RUSTFS_RUNTIME_FSYNC_BLOCKING_THREADS";
|
||||||
|
pub const DEFAULT_FSYNC_BLOCKING_THREADS: usize = 0;
|
||||||
|
|
||||||
// Dial9 Tokio Telemetry Default values
|
// Dial9 Tokio Telemetry Default values
|
||||||
pub const DEFAULT_RUNTIME_DIAL9_ENABLED: bool = false; // Disabled by default
|
pub const DEFAULT_RUNTIME_DIAL9_ENABLED: bool = false; // Disabled by default
|
||||||
pub const DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR: &str = "/var/log/rustfs/telemetry";
|
pub const DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR: &str = "/var/log/rustfs/telemetry";
|
||||||
|
|||||||
@@ -3275,9 +3275,6 @@ impl ECStore {
|
|||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
let source_cleanup_mutation_fence = self
|
|
||||||
.acquire_decommission_source_cleanup_fence(bucket.as_str(), entry.name.as_str(), set.as_ref())
|
|
||||||
.await?;
|
|
||||||
let cleanup_result = data_movement::cleanup_source_entry_if_unchanged(
|
let cleanup_result = data_movement::cleanup_source_entry_if_unchanged(
|
||||||
set.clone(),
|
set.clone(),
|
||||||
bucket.as_str(),
|
bucket.as_str(),
|
||||||
@@ -3289,7 +3286,6 @@ impl ECStore {
|
|||||||
lifecycle_guard: bucket_incarnation_fence
|
lifecycle_guard: bucket_incarnation_fence
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.and_then(|guard| guard.namespace_lock_guard()),
|
.and_then(|guard| guard.namespace_lock_guard()),
|
||||||
object_mutation_fence: Some(&source_cleanup_mutation_fence),
|
|
||||||
},
|
},
|
||||||
"decommission",
|
"decommission",
|
||||||
)
|
)
|
||||||
@@ -3383,22 +3379,6 @@ impl ECStore {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) async fn decommission_entry_for_test(
|
|
||||||
self: &Arc<Self>,
|
|
||||||
idx: usize,
|
|
||||||
entry: MetaCacheEntry,
|
|
||||||
bucket: String,
|
|
||||||
set: Arc<SetDisks>,
|
|
||||||
) -> Result<()> {
|
|
||||||
let worker_permit = Arc::new(Semaphore::new(1))
|
|
||||||
.acquire_owned()
|
|
||||||
.await
|
|
||||||
.map_err(|err| Error::other(format!("decommission test worker permit acquire failed: {err}")))?;
|
|
||||||
self.decommission_entry(CancellationToken::new(), idx, entry, bucket, set, worker_permit, None, None, None, None)
|
|
||||||
.await
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tracing::instrument(skip(self, rx))]
|
#[tracing::instrument(skip(self, rx))]
|
||||||
async fn decommission_pool(
|
async fn decommission_pool(
|
||||||
self: &Arc<Self>,
|
self: &Arc<Self>,
|
||||||
@@ -4334,20 +4314,15 @@ impl ECStore {
|
|||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name);
|
warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name);
|
||||||
let object_name = rd.object_info.name.clone();
|
let object_name = rd.object_info.name.clone();
|
||||||
let mut migration = tokio::task::JoinSet::new();
|
let result = data_movement::migrate_object(
|
||||||
migration.spawn(data_movement::migrate_decommission_object(
|
|
||||||
self,
|
self,
|
||||||
pool_idx,
|
pool_idx,
|
||||||
bucket.clone(),
|
bucket.clone(),
|
||||||
rd,
|
rd,
|
||||||
expected_bucket_incarnation_id,
|
expected_bucket_incarnation_id,
|
||||||
"decommission_object",
|
"decommission_object",
|
||||||
));
|
)
|
||||||
let result = migration
|
.await;
|
||||||
.join_next()
|
|
||||||
.await
|
|
||||||
.ok_or_else(|| Error::other("decommission migration task was not started"))?
|
|
||||||
.map_err(|err| Error::other(format!("decommission migration task join error: {err}")))?;
|
|
||||||
if result.is_ok() {
|
if result.is_ok() {
|
||||||
warn!("decommission_object: migrated {} {}", &bucket, &object_name);
|
warn!("decommission_object: migrated {} {}", &bucket, &object_name);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ use crate::storage_api_contracts::{
|
|||||||
namespace::NamespaceLocking as _,
|
namespace::NamespaceLocking as _,
|
||||||
object::{HTTPPreconditions, ObjectOperations as _},
|
object::{HTTPPreconditions, ObjectOperations as _},
|
||||||
};
|
};
|
||||||
use crate::store::{ECStore, ObjectLockDiagGuard, SourceCleanupMutationFence};
|
use crate::store::ECStore;
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo};
|
use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo};
|
||||||
use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex};
|
use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex};
|
||||||
@@ -856,6 +856,7 @@ fn is_equivalent_data_movement_object(source: &ObjectInfo, target: &ObjectInfo)
|
|||||||
fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool {
|
fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool {
|
||||||
is_unversioned_data_movement_object(source)
|
is_unversioned_data_movement_object(source)
|
||||||
&& is_unversioned_data_movement_object(target)
|
&& is_unversioned_data_movement_object(target)
|
||||||
|
&& !target.delete_marker
|
||||||
&& source
|
&& source
|
||||||
.mod_time
|
.mod_time
|
||||||
.zip(target.mod_time)
|
.zip(target.mod_time)
|
||||||
@@ -1027,7 +1028,6 @@ pub(crate) enum SourceCleanupError {
|
|||||||
pub(crate) struct SourceCleanupBucketFence<'a> {
|
pub(crate) struct SourceCleanupBucketFence<'a> {
|
||||||
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
|
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
|
||||||
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
|
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
|
||||||
pub(crate) object_mutation_fence: Option<&'a SourceCleanupMutationFence>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn ensure_source_cleanup_versions_match(
|
fn ensure_source_cleanup_versions_match(
|
||||||
@@ -1065,9 +1065,7 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
|
|||||||
struct SourceCleanupDeleteBarrierState {
|
struct SourceCleanupDeleteBarrierState {
|
||||||
bucket: String,
|
bucket: String,
|
||||||
object: String,
|
object: String,
|
||||||
fence_pending: tokio::sync::Notify,
|
|
||||||
arrived: tokio::sync::Notify,
|
arrived: tokio::sync::Notify,
|
||||||
is_paused: AtomicBool,
|
|
||||||
release: tokio::sync::Notify,
|
release: tokio::sync::Notify,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1081,7 +1079,7 @@ pub(crate) struct SourceCleanupDeleteBarrier {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<SourceCleanupDeleteBarrierState>>>> =
|
static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<SourceCleanupDeleteBarrierState>>>> =
|
||||||
std::sync::OnceLock::new();
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -1094,22 +1092,15 @@ impl SourceCleanupDeleteBarrier {
|
|||||||
let state = Arc::new(SourceCleanupDeleteBarrierState {
|
let state = Arc::new(SourceCleanupDeleteBarrierState {
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
object: object.to_string(),
|
object: object.to_string(),
|
||||||
fence_pending: tokio::sync::Notify::new(),
|
|
||||||
arrived: tokio::sync::Notify::new(),
|
arrived: tokio::sync::Notify::new(),
|
||||||
is_paused: AtomicBool::new(false),
|
|
||||||
release: tokio::sync::Notify::new(),
|
release: tokio::sync::Notify::new(),
|
||||||
});
|
});
|
||||||
let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS
|
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
.lock()
|
.lock()
|
||||||
.expect("source cleanup delete barrier mutex should not poison");
|
.expect("source cleanup delete barrier mutex should not poison");
|
||||||
assert!(
|
assert!(slot.is_none(), "source cleanup delete barrier must be unique");
|
||||||
!barriers
|
*slot = Some(Arc::clone(&state));
|
||||||
.iter()
|
|
||||||
.any(|barrier| barrier.bucket == bucket && barrier.object == object),
|
|
||||||
"source cleanup delete barrier must be unique per object"
|
|
||||||
);
|
|
||||||
barriers.push(Arc::clone(&state));
|
|
||||||
Self { state }
|
Self { state }
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1119,58 +1110,35 @@ impl SourceCleanupDeleteBarrier {
|
|||||||
.expect("source cleanup should reach the pre-delete barrier");
|
.expect("source cleanup should reach the pre-delete barrier");
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn wait_until_fence_pending(&self) {
|
|
||||||
tokio::time::timeout(StdDuration::from_secs(30), self.state.fence_pending.notified())
|
|
||||||
.await
|
|
||||||
.expect("source cleanup should attempt the fixed mutation fence");
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn is_paused(&self) -> bool {
|
|
||||||
self.state.is_paused.load(Ordering::Acquire)
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn release(&self) {
|
pub(crate) fn release(&self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) fn notify_source_cleanup_mutation_fence_pending(bucket: &str, object: &str) {
|
|
||||||
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
|
||||||
.lock()
|
|
||||||
.expect("source cleanup delete barrier mutex should not poison")
|
|
||||||
.iter()
|
|
||||||
.find(|barrier| barrier.bucket == bucket && barrier.object == object)
|
|
||||||
.cloned();
|
|
||||||
if let Some(barrier) = barrier {
|
|
||||||
barrier.fence_pending.notify_one();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
impl Drop for SourceCleanupDeleteBarrier {
|
impl Drop for SourceCleanupDeleteBarrier {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS
|
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
.lock()
|
.lock()
|
||||||
.expect("source cleanup delete barrier mutex should not poison");
|
.expect("source cleanup delete barrier mutex should not poison");
|
||||||
barriers.retain(|state| !Arc::ptr_eq(state, &self.state));
|
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
||||||
|
*slot = None;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
|
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
|
||||||
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
|
let barrier = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
.lock()
|
.lock()
|
||||||
.expect("source cleanup delete barrier mutex should not poison")
|
.expect("source cleanup delete barrier mutex should not poison")
|
||||||
.iter()
|
.as_ref()
|
||||||
.find(|barrier| barrier.bucket == bucket && barrier.object == object)
|
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
|
||||||
.cloned();
|
.cloned();
|
||||||
if let Some(barrier) = barrier {
|
if let Some(barrier) = barrier {
|
||||||
barrier.is_paused.store(true, Ordering::Release);
|
|
||||||
barrier.arrived.notify_one();
|
barrier.arrived.notify_one();
|
||||||
barrier.release.notified().await;
|
barrier.release.notified().await;
|
||||||
}
|
}
|
||||||
@@ -1186,20 +1154,11 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
op_label: &str,
|
op_label: &str,
|
||||||
) -> std::result::Result<ObjectInfo, SourceCleanupError> {
|
) -> std::result::Result<ObjectInfo, SourceCleanupError> {
|
||||||
let cleanup_key = encode_dir_object(object);
|
let cleanup_key = encode_dir_object(object);
|
||||||
let source_guard = if bucket_fence
|
let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?;
|
||||||
.object_mutation_fence
|
let _guard = ns_lock
|
||||||
.is_some_and(SourceCleanupMutationFence::source_lock_covered)
|
.get_write_lock(get_lock_acquire_timeout())
|
||||||
{
|
.await
|
||||||
None
|
.map_err(Error::from)?;
|
||||||
} else {
|
|
||||||
let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?;
|
|
||||||
Some(
|
|
||||||
ns_lock
|
|
||||||
.get_write_lock(get_lock_acquire_timeout())
|
|
||||||
.await
|
|
||||||
.map_err(Error::from)?,
|
|
||||||
)
|
|
||||||
};
|
|
||||||
|
|
||||||
if bucket_fence
|
if bucket_fence
|
||||||
.lifecycle_guard
|
.lifecycle_guard
|
||||||
@@ -1209,14 +1168,6 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
"{op_label}: bucket incarnation fence was lost before source cleanup"
|
"{op_label}: bucket incarnation fence was lost before source cleanup"
|
||||||
))));
|
))));
|
||||||
}
|
}
|
||||||
if bucket_fence
|
|
||||||
.object_mutation_fence
|
|
||||||
.is_some_and(SourceCleanupMutationFence::is_lock_lost)
|
|
||||||
{
|
|
||||||
return Err(SourceCleanupError::Storage(Error::other(format!(
|
|
||||||
"{op_label}: object mutation fence was lost before source cleanup"
|
|
||||||
))));
|
|
||||||
}
|
|
||||||
|
|
||||||
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
|
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
|
||||||
|
|
||||||
@@ -1231,12 +1182,7 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
|
|||||||
expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id,
|
expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
if let Some(source_guard) = source_guard.as_ref() {
|
opts.add_namespace_lock_guard(&_guard);
|
||||||
opts.add_namespace_lock_guard(source_guard);
|
|
||||||
}
|
|
||||||
if let Some(object_mutation_fence) = bucket_fence.object_mutation_fence {
|
|
||||||
object_mutation_fence.add_namespace_lock_fence(&mut opts);
|
|
||||||
}
|
|
||||||
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
|
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
|
||||||
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
|
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
|
||||||
}
|
}
|
||||||
@@ -1384,37 +1330,6 @@ fn data_movement_part_upload_failure_stage(err: &Error) -> &'static str {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn migrate_decommission_object(
|
|
||||||
store: Arc<ECStore>,
|
|
||||||
pool_idx: usize,
|
|
||||||
bucket: String,
|
|
||||||
rd: GetObjectReader,
|
|
||||||
source_bucket_incarnation_id: Option<uuid::Uuid>,
|
|
||||||
op_label: &str,
|
|
||||||
) -> Result<()> {
|
|
||||||
let source = rd.object_info.clone();
|
|
||||||
let _mutation_fence = store
|
|
||||||
.acquire_decommission_object_mutation_fence(&bucket, &source.name)
|
|
||||||
.await?;
|
|
||||||
let current = find_data_movement_target_info(store.as_ref(), pool_idx, &bucket, &source)
|
|
||||||
.await?
|
|
||||||
.ok_or(Error::FileNotFound)?;
|
|
||||||
if !is_equivalent_data_movement_object_identity(&source, ¤t, true, false) {
|
|
||||||
return Err(Error::FileNotFound);
|
|
||||||
}
|
|
||||||
|
|
||||||
migrate_object_inner(
|
|
||||||
store,
|
|
||||||
pool_idx,
|
|
||||||
bucket,
|
|
||||||
rd,
|
|
||||||
source_bucket_incarnation_id,
|
|
||||||
op_label,
|
|
||||||
Some(&_mutation_fence),
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) async fn migrate_object(
|
pub(crate) async fn migrate_object(
|
||||||
store: Arc<ECStore>,
|
store: Arc<ECStore>,
|
||||||
pool_idx: usize,
|
pool_idx: usize,
|
||||||
@@ -1422,18 +1337,6 @@ pub(crate) async fn migrate_object(
|
|||||||
rd: GetObjectReader,
|
rd: GetObjectReader,
|
||||||
source_bucket_incarnation_id: Option<uuid::Uuid>,
|
source_bucket_incarnation_id: Option<uuid::Uuid>,
|
||||||
op_label: &str,
|
op_label: &str,
|
||||||
) -> Result<()> {
|
|
||||||
migrate_object_inner(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn migrate_object_inner(
|
|
||||||
store: Arc<ECStore>,
|
|
||||||
pool_idx: usize,
|
|
||||||
bucket: String,
|
|
||||||
rd: GetObjectReader,
|
|
||||||
source_bucket_incarnation_id: Option<uuid::Uuid>,
|
|
||||||
op_label: &str,
|
|
||||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let object_info = rd.object_info.clone();
|
let object_info = rd.object_info.clone();
|
||||||
let has_part_checksums = object_info
|
let has_part_checksums = object_info
|
||||||
@@ -1447,7 +1350,7 @@ async fn migrate_object_inner(
|
|||||||
let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx);
|
let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx);
|
||||||
new_multipart_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
|
new_multipart_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
|
||||||
let (res, target_pool_idx, expected_bucket_incarnation_id) = match store
|
let (res, target_pool_idx, expected_bucket_incarnation_id) = match store
|
||||||
.handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts, mutation_fence)
|
.handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(res) => res,
|
Ok(res) => res,
|
||||||
@@ -1545,7 +1448,7 @@ async fn migrate_object_inner(
|
|||||||
if let Err(err) = store
|
if let Err(err) = store
|
||||||
.clone()
|
.clone()
|
||||||
.complete_multipart_upload_for_data_movement(
|
.complete_multipart_upload_for_data_movement(
|
||||||
(target_pool_idx, mutation_fence),
|
target_pool_idx,
|
||||||
&bucket,
|
&bucket,
|
||||||
&object_info.name,
|
&object_info.name,
|
||||||
&res.upload_id,
|
&res.upload_id,
|
||||||
@@ -1706,7 +1609,7 @@ async fn migrate_object_inner(
|
|||||||
let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx);
|
let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx);
|
||||||
put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
|
put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
|
||||||
let (target_pool_idx, put_result) = store
|
let (target_pool_idx, put_result) = store
|
||||||
.put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts, mutation_fence)
|
.put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| data_movement_stage_error(op_label, "prepare_put_object", &bucket, &object_info.name, err))?;
|
.map_err(|err| data_movement_stage_error(op_label, "prepare_put_object", &bucket, &object_info.name, err))?;
|
||||||
if let Err(err) = put_result {
|
if let Err(err) = put_result {
|
||||||
@@ -3638,47 +3541,25 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_precondition_conflict_accepts_only_newer_null_delete_marker() {
|
fn test_precondition_conflict_rejects_newer_delete_marker() {
|
||||||
for version_id in [None, Some(Uuid::nil())] {
|
let source = ObjectInfo {
|
||||||
let source = ObjectInfo {
|
size: 128,
|
||||||
version_id,
|
etag: Some("etag-source".to_string()),
|
||||||
size: 128,
|
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
|
||||||
etag: Some("etag-source".to_string()),
|
..Default::default()
|
||||||
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
|
};
|
||||||
..Default::default()
|
let target = ObjectInfo {
|
||||||
};
|
delete_marker: true,
|
||||||
let target = ObjectInfo {
|
etag: None,
|
||||||
delete_marker: true,
|
mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND),
|
||||||
etag: None,
|
..source.clone()
|
||||||
mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND),
|
};
|
||||||
..source.clone()
|
|
||||||
};
|
|
||||||
|
|
||||||
assert!(
|
let should_resume =
|
||||||
resolve_data_movement_overwrite_resume_result(
|
resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(target)), &source, 0, 1)
|
||||||
&Error::PreconditionFailed,
|
.expect("delete marker conflict should be evaluated");
|
||||||
Ok(Some(target.clone())),
|
|
||||||
&source,
|
|
||||||
0,
|
|
||||||
1,
|
|
||||||
)
|
|
||||||
.expect("newer null delete marker should be evaluated")
|
|
||||||
);
|
|
||||||
|
|
||||||
let mut same_time = target.clone();
|
assert!(!should_resume);
|
||||||
same_time.mod_time = source.mod_time;
|
|
||||||
assert!(
|
|
||||||
!resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(same_time)), &source, 0, 1,)
|
|
||||||
.expect("same-generation null delete marker should be rejected")
|
|
||||||
);
|
|
||||||
|
|
||||||
let mut versioned = target;
|
|
||||||
versioned.version_id = Some(Uuid::new_v4());
|
|
||||||
assert!(
|
|
||||||
!resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(versioned)), &source, 0, 1,)
|
|
||||||
.expect("a UUID delete marker must not erase a null source version")
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
@@ -858,7 +858,6 @@ const EVENT_DISK_LOCAL_DIRECT_IO_FALLBACK: &str = "disk_local_direct_io_fallback
|
|||||||
#[cfg(target_os = "linux")]
|
#[cfg(target_os = "linux")]
|
||||||
const EVENT_DISK_LOCAL_URING_LATCH_OFF: &str = "disk_local_uring_latch_off";
|
const EVENT_DISK_LOCAL_URING_LATCH_OFF: &str = "disk_local_uring_latch_off";
|
||||||
const EVENT_DISK_LOCAL_DELETE_FAILED: &str = "disk_local_delete_failed";
|
const EVENT_DISK_LOCAL_DELETE_FAILED: &str = "disk_local_delete_failed";
|
||||||
const EVENT_DISK_LOCAL_DELETE_ROLLBACK_FAILED: &str = "disk_local_delete_rollback_failed";
|
|
||||||
const EVENT_DISK_LOCAL_CHECK_PARTS: &str = "disk_local_check_parts";
|
const EVENT_DISK_LOCAL_CHECK_PARTS: &str = "disk_local_check_parts";
|
||||||
const EVENT_DISK_LOCAL_ACCESS_FAILED: &str = "disk_local_access_failed";
|
const EVENT_DISK_LOCAL_ACCESS_FAILED: &str = "disk_local_access_failed";
|
||||||
const EVENT_DISK_LOCAL_VOLUME_SETUP_FAILED: &str = "disk_local_volume_setup_failed";
|
const EVENT_DISK_LOCAL_VOLUME_SETUP_FAILED: &str = "disk_local_volume_setup_failed";
|
||||||
@@ -6107,43 +6106,6 @@ impl LocalDisk {
|
|||||||
Ok((bytes, modtime))
|
Ok((bytes, modtime))
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn write_missing_delete_marker(
|
|
||||||
&self,
|
|
||||||
volume: &str,
|
|
||||||
path: &str,
|
|
||||||
fi: FileInfo,
|
|
||||||
object_dir: &Path,
|
|
||||||
xl_path: &Path,
|
|
||||||
rollback_dir: Option<Uuid>,
|
|
||||||
) -> Result<()> {
|
|
||||||
if let Some(rollback_dir) = rollback_dir {
|
|
||||||
let rollback_path = object_dir.join(rollback_dir.to_string());
|
|
||||||
fs::create_dir_all(&rollback_path).await.map_err(to_file_error)?;
|
|
||||||
fs::write(rollback_path.join(DELETE_MARKER_ROLLBACK_FILE), [])
|
|
||||||
.await
|
|
||||||
.map_err(to_file_error)?;
|
|
||||||
}
|
|
||||||
if let Err(err) = self.write_metadata("", volume, path, fi).await {
|
|
||||||
if let Some(rollback_dir) = rollback_dir
|
|
||||||
&& let Err(restore_err) = restore_delete_rollback(object_dir, xl_path, rollback_dir, &self.publication_root).await
|
|
||||||
{
|
|
||||||
warn!(
|
|
||||||
event = EVENT_DISK_LOCAL_DELETE_ROLLBACK_FAILED,
|
|
||||||
component = LOG_COMPONENT_ECSTORE,
|
|
||||||
subsystem = LOG_SUBSYSTEM_DISK_LOCAL,
|
|
||||||
result = "failed",
|
|
||||||
volume,
|
|
||||||
path,
|
|
||||||
rollback_dir = %rollback_dir,
|
|
||||||
error = ?restore_err,
|
|
||||||
"Disk local delete rollback failed"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
return Err(err);
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &[FileInfo], opts: &DeleteOptions) -> Result<()> {
|
async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &[FileInfo], opts: &DeleteOptions) -> Result<()> {
|
||||||
let volume_dir = self.io_get_bucket_path(volume)?;
|
let volume_dir = self.io_get_bucket_path(volume)?;
|
||||||
let xlpath = self.io_get_object_path(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str())?;
|
let xlpath = self.io_get_object_path(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str())?;
|
||||||
@@ -6161,20 +6123,7 @@ impl LocalDisk {
|
|||||||
return restore_metadata_backup(object_dir, &xlpath, rollback_dir, &self.publication_root).await;
|
return restore_metadata_backup(object_dir, &xlpath, rollback_dir, &self.publication_root).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
let (data, _) = match self.read_all_data_with_dmtime(volume, volume_dir.as_path(), &xlpath).await {
|
let (data, _) = self.read_all_data_with_dmtime(volume, volume_dir.as_path(), &xlpath).await?;
|
||||||
Ok(data) => data,
|
|
||||||
Err(DiskError::FileNotFound) => {
|
|
||||||
// `deleted` alone can be an explicit marker purge; only
|
|
||||||
// `mark_deleted` may create metadata that was not present.
|
|
||||||
let Some(delete_marker) = fis.iter().find(|fi| fi.deleted && fi.mark_deleted).cloned() else {
|
|
||||||
return Err(DiskError::FileNotFound);
|
|
||||||
};
|
|
||||||
return self
|
|
||||||
.write_missing_delete_marker(volume, path, delete_marker, object_dir, &xlpath, opts.old_data_dir)
|
|
||||||
.await;
|
|
||||||
}
|
|
||||||
Err(err) => return Err(err),
|
|
||||||
};
|
|
||||||
|
|
||||||
if data.is_empty() {
|
if data.is_empty() {
|
||||||
return Err(DiskError::FileNotFound);
|
return Err(DiskError::FileNotFound);
|
||||||
@@ -10473,9 +10422,29 @@ impl DiskAPI for LocalDisk {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if fi.deleted && force_del_marker {
|
if fi.deleted && force_del_marker {
|
||||||
return self
|
if let Some(rollback_dir) = rollback_dir {
|
||||||
.write_missing_delete_marker(volume, path, fi, file_path.as_path(), &xl_path, rollback_dir)
|
let rollback_path = file_path.join(rollback_dir.to_string());
|
||||||
.await;
|
fs::create_dir_all(&rollback_path).await.map_err(to_file_error)?;
|
||||||
|
fs::write(rollback_path.join(DELETE_MARKER_ROLLBACK_FILE), [])
|
||||||
|
.await
|
||||||
|
.map_err(to_file_error)?;
|
||||||
|
}
|
||||||
|
if let Err(err) = self.write_metadata("", volume, path, fi).await {
|
||||||
|
if let Some(rollback_dir) = rollback_dir
|
||||||
|
&& let Err(restore_err) =
|
||||||
|
restore_delete_rollback(file_path.as_path(), &xl_path, rollback_dir, &self.publication_root).await
|
||||||
|
{
|
||||||
|
warn!(
|
||||||
|
volume,
|
||||||
|
path,
|
||||||
|
rollback_dir = %rollback_dir,
|
||||||
|
error = ?restore_err,
|
||||||
|
"failed to restore metadata after delete marker commit error"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return Err(err);
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
return if fi.version_id.is_some() {
|
return if fi.version_id.is_some() {
|
||||||
|
|||||||
@@ -315,7 +315,7 @@ pub async fn fsync_dir(dir: impl AsRef<Path>) -> io::Result<()> {
|
|||||||
#[cfg(unix)]
|
#[cfg(unix)]
|
||||||
{
|
{
|
||||||
let dir = dir.as_ref().to_path_buf();
|
let dir = dir.as_ref().to_path_buf();
|
||||||
tokio::task::spawn_blocking(move || fsync_dir_std(dir)).await?
|
fsync_spawn_blocking(move || fsync_dir_std(dir)).await?
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(not(unix))]
|
#[cfg(not(unix))]
|
||||||
@@ -683,7 +683,7 @@ async fn fsync_open_dst_dir_group(group: &DstDirFsyncGroup) -> io::Result<()> {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
let dir = group.dir.clone();
|
let dir = group.dir.clone();
|
||||||
let dir_file = group.dir_file.clone();
|
let dir_file = group.dir_file.clone();
|
||||||
tokio::task::spawn_blocking(move || {
|
fsync_spawn_blocking(move || {
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
{
|
{
|
||||||
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
|
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
|
||||||
@@ -1080,6 +1080,44 @@ const TEST_GLOBAL_FILE_SYNCS: usize = 64;
|
|||||||
|
|
||||||
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
||||||
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||||
|
|
||||||
|
/// Dedicated tokio runtime for fsync/fdatasync blocking operations. When
|
||||||
|
/// configured with >1 threads, isolates device-bound fsync from the main
|
||||||
|
/// blocking pool so reads (pread/stat/open) are not starved. `None` means
|
||||||
|
/// fall back to the main runtime (zero behavior change).
|
||||||
|
static FSYNC_RUNTIME: LazyLock<Option<tokio::runtime::Runtime>> = LazyLock::new(|| {
|
||||||
|
let threads =
|
||||||
|
rustfs_utils::get_env_usize(rustfs_config::ENV_FSYNC_BLOCKING_THREADS, rustfs_config::DEFAULT_FSYNC_BLOCKING_THREADS);
|
||||||
|
if threads <= 1 {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let mut builder = tokio::runtime::Builder::new_multi_thread();
|
||||||
|
builder
|
||||||
|
.worker_threads(num_cpus::get().min(8))
|
||||||
|
.max_blocking_threads(threads)
|
||||||
|
.thread_name("rustfs-fsync")
|
||||||
|
.thread_stack_size(512 * 1024)
|
||||||
|
.enable_all();
|
||||||
|
match builder.build() {
|
||||||
|
Ok(rt) => {
|
||||||
|
tracing::info!(threads, "fsync dedicated blocking pool enabled");
|
||||||
|
Some(rt)
|
||||||
|
}
|
||||||
|
Err(err) => {
|
||||||
|
tracing::warn!(%err, "failed to build fsync runtime, falling back to main pool");
|
||||||
|
None
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
/// Spawn a blocking task on the fsync-dedicated runtime if configured,
|
||||||
|
/// otherwise fall back to the main tokio blocking pool.
|
||||||
|
fn fsync_spawn_blocking<T: Send + 'static>(f: impl FnOnce() -> T + Send + 'static) -> tokio::task::JoinHandle<T> {
|
||||||
|
match FSYNC_RUNTIME.as_ref() {
|
||||||
|
Some(rt) => rt.spawn_blocking(f),
|
||||||
|
None => tokio::task::spawn_blocking(f),
|
||||||
|
}
|
||||||
|
}
|
||||||
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
||||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||||
type NamespaceMutationLock = AsyncMutex<()>;
|
type NamespaceMutationLock = AsyncMutex<()>;
|
||||||
@@ -1217,7 +1255,7 @@ where
|
|||||||
F: FnOnce() -> io::Result<T> + Send + 'static,
|
F: FnOnce() -> io::Result<T> + Send + 'static,
|
||||||
{
|
{
|
||||||
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
|
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
|
||||||
let result = tokio::task::spawn_blocking(move || {
|
let result = fsync_spawn_blocking(move || {
|
||||||
let _disk_permit = disk_permit;
|
let _disk_permit = disk_permit;
|
||||||
work()
|
work()
|
||||||
})
|
})
|
||||||
@@ -2146,7 +2184,7 @@ async fn run_blocking_namespace_file_sync_operation_with_global<T: Send + 'stati
|
|||||||
wait_started,
|
wait_started,
|
||||||
);
|
);
|
||||||
let disk_permit = admission.disk_permit.clone();
|
let disk_permit = admission.disk_permit.clone();
|
||||||
let result = tokio::task::spawn_blocking(move || {
|
let result = fsync_spawn_blocking(move || {
|
||||||
let _lease = lease;
|
let _lease = lease;
|
||||||
let _disk_permit = disk_permit;
|
let _disk_permit = disk_permit;
|
||||||
operation()
|
operation()
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ use crate::storage_api_contracts::{
|
|||||||
pub struct NamespaceLockFence {
|
pub struct NamespaceLockFence {
|
||||||
signals: Arc<Vec<Arc<rustfs_lock::distributed_lock::LockLostSignal>>>,
|
signals: Arc<Vec<Arc<rustfs_lock::distributed_lock::LockLostSignal>>>,
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
forced_lost: Arc<Vec<Arc<std::sync::atomic::AtomicBool>>>,
|
forced_lost: Arc<std::sync::atomic::AtomicBool>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Debug for NamespaceLockFence {
|
impl Debug for NamespaceLockFence {
|
||||||
@@ -40,17 +40,13 @@ impl NamespaceLockFence {
|
|||||||
Self {
|
Self {
|
||||||
signals: Arc::default(),
|
signals: Arc::default(),
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
forced_lost: Arc::new(vec![Arc::new(std::sync::atomic::AtomicBool::new(false))]),
|
forced_lost: Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn is_lock_lost(&self) -> bool {
|
pub(crate) fn is_lock_lost(&self) -> bool {
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
if self
|
if self.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
|
||||||
.forced_lost
|
|
||||||
.iter()
|
|
||||||
.any(|lost| lost.load(std::sync::atomic::Ordering::Acquire))
|
|
||||||
{
|
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
self.signals.iter().any(|signal| signal.is_lost())
|
self.signals.iter().any(|signal| signal.is_lost())
|
||||||
@@ -61,26 +57,27 @@ impl NamespaceLockFence {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn extend(&mut self, other: &Self) {
|
fn extend(&mut self, other: &Self) {
|
||||||
if !Arc::ptr_eq(&self.signals, &other.signals) {
|
if Arc::ptr_eq(&self.signals, &other.signals) {
|
||||||
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
|
return;
|
||||||
}
|
}
|
||||||
|
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
if !Arc::ptr_eq(&self.forced_lost, &other.forced_lost) {
|
if other.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
|
||||||
Arc::make_mut(&mut self.forced_lost).extend(other.forced_lost.iter().cloned());
|
self.forced_lost.store(true, std::sync::atomic::Ordering::Release);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pub(crate) fn lost_for_test() -> Self {
|
pub(crate) fn lost_for_test() -> Self {
|
||||||
let fence = Self::new();
|
let fence = Self::new();
|
||||||
fence.forced_lost[0].store(true, std::sync::atomic::Ordering::Release);
|
fence.forced_lost.store(true, std::sync::atomic::Ordering::Release);
|
||||||
fence
|
fence
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pub(crate) fn loss_handle_for_test() -> (Self, Arc<std::sync::atomic::AtomicBool>) {
|
pub(crate) fn loss_handle_for_test() -> (Self, Arc<std::sync::atomic::AtomicBool>) {
|
||||||
let fence = Self::new();
|
let fence = Self::new();
|
||||||
(fence.clone(), Arc::clone(&fence.forced_lost[0]))
|
(fence.clone(), Arc::clone(&fence.forced_lost))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -414,13 +411,6 @@ impl ObjectOptions {
|
|||||||
self.namespace_lock_fence.get_or_insert_with(NamespaceLockFence::new);
|
self.namespace_lock_fence.get_or_insert_with(NamespaceLockFence::new);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) fn add_namespace_lock_fence_for_test(&mut self, fence: &NamespaceLockFence) {
|
|
||||||
self.namespace_lock_fence
|
|
||||||
.get_or_insert_with(NamespaceLockFence::new)
|
|
||||||
.extend(fence);
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn ensure_lifecycle_delete_all_journal(&mut self) {
|
pub(crate) fn ensure_lifecycle_delete_all_journal(&mut self) {
|
||||||
self.lifecycle_delete_all_journal
|
self.lifecycle_delete_all_journal
|
||||||
.get_or_insert_with(|| Arc::new(parking_lot::Mutex::new(LifecycleDeleteAllJournalState::default())));
|
.get_or_insert_with(|| Arc::new(parking_lot::Mutex::new(LifecycleDeleteAllJournalState::default())));
|
||||||
|
|||||||
@@ -334,7 +334,6 @@ impl ECStore {
|
|||||||
lifecycle_guard: bucket_incarnation_fence
|
lifecycle_guard: bucket_incarnation_fence
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.and_then(|guard| guard.namespace_lock_guard()),
|
.and_then(|guard| guard.namespace_lock_guard()),
|
||||||
..Default::default()
|
|
||||||
},
|
},
|
||||||
"rebalance",
|
"rebalance",
|
||||||
),
|
),
|
||||||
|
|||||||
@@ -735,12 +735,8 @@ pub(crate) use core::io_primitives::disk_call_counters;
|
|||||||
mod ctx;
|
mod ctx;
|
||||||
mod metadata;
|
mod metadata;
|
||||||
mod ops;
|
mod ops;
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) use ops::multipart::NewMultipartUploadCommitObservation;
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
|
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) use ops::object::DeleteObjectCommitBarrier;
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
||||||
pub(crate) use ops::object::body_cache_plaintext_len;
|
pub(crate) use ops::object::body_cache_plaintext_len;
|
||||||
@@ -3029,16 +3025,6 @@ pub struct SetDisks {
|
|||||||
storage_class_config_override: Arc<std::sync::RwLock<Option<Arc<storageclass::Config>>>>,
|
storage_class_config_override: Arc<std::sync::RwLock<Option<Arc<storageclass::Config>>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
// DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones
|
|
||||||
// each endpoint's canonical Arc, so an exact Arc set identifies the lock domain.
|
|
||||||
pub(crate) fn same_distributed_lock_domain(left: &[Arc<dyn LockClient>], right: &[Arc<dyn LockClient>]) -> bool {
|
|
||||||
left.iter()
|
|
||||||
.all(|left_client| right.iter().any(|right_client| Arc::ptr_eq(left_client, right_client)))
|
|
||||||
&& right
|
|
||||||
.iter()
|
|
||||||
.all(|right_client| left.iter().any(|left_client| Arc::ptr_eq(left_client, right_client)))
|
|
||||||
}
|
|
||||||
|
|
||||||
const ERASURE_CACHE_MAX_ENTRIES: usize = 32;
|
const ERASURE_CACHE_MAX_ENTRIES: usize = 32;
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||||
@@ -3614,15 +3600,6 @@ impl SetDisks {
|
|||||||
&self.ctx
|
&self.ctx
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Whether both sets' namespace-lock implementations cover the same object key.
|
|
||||||
pub(crate) async fn shares_namespace_lock_domain(&self, other: &Self) -> bool {
|
|
||||||
match (self.ctx.is_dist_erasure().await, other.ctx.is_dist_erasure().await) {
|
|
||||||
(false, false) => Arc::ptr_eq(&self.local_lock_manager, &other.local_lock_manager),
|
|
||||||
(true, true) => same_distributed_lock_domain(&self.lockers, &other.lockers),
|
|
||||||
_ => false,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The lock manager this set actually uses (test-only; Phase 5 Slice 3).
|
/// The lock manager this set actually uses (test-only; Phase 5 Slice 3).
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pub(crate) fn local_lock_manager_for_test(&self) -> &Arc<rustfs_lock::GlobalLockManager> {
|
pub(crate) fn local_lock_manager_for_test(&self) -> &Arc<rustfs_lock::GlobalLockManager> {
|
||||||
@@ -4607,11 +4584,11 @@ fn should_preserve_delete_replication_state(opts: &ObjectOptions) -> bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn should_force_delete_marker_for_missing_version(opts: &ObjectOptions) -> bool {
|
fn should_force_delete_marker_for_missing_version(opts: &ObjectOptions) -> bool {
|
||||||
opts.delete_marker || ((opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.data_movement)
|
opts.delete_marker || (opts.versioned && opts.version_id.is_none() && !opts.data_movement)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_found: bool) -> (bool, bool) {
|
fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_found: bool) -> (bool, bool) {
|
||||||
let mut mark_delete = goi.version_id.is_some() || ((opts.versioned || opts.version_suspended) && opts.version_id.is_none());
|
let mut mark_delete = goi.version_id.is_some() || (opts.versioned && opts.version_id.is_none());
|
||||||
let mut delete_marker = opts.versioned;
|
let mut delete_marker = opts.versioned;
|
||||||
|
|
||||||
if opts.version_id.is_some() {
|
if opts.version_id.is_some() {
|
||||||
|
|||||||
@@ -32,8 +32,6 @@ use crate::crash_inject::{self, CrashPoint};
|
|||||||
use crate::multipart_listing::paginate_multipart_listing;
|
use crate::multipart_listing::paginate_multipart_listing;
|
||||||
use futures::{StreamExt, stream};
|
use futures::{StreamExt, stream};
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
#[cfg(test)]
|
|
||||||
use std::sync::atomic::AtomicBool;
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
@@ -67,7 +65,6 @@ impl StaleMultipartCleanupGuard {
|
|||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
#[derive(Clone, Copy, PartialEq, Eq)]
|
||||||
pub enum MultipartCommitPause {
|
pub enum MultipartCommitPause {
|
||||||
NewUploadBeforeLockLost,
|
|
||||||
PutPartBeforeLockAcquire,
|
PutPartBeforeLockAcquire,
|
||||||
PutPartBeforeLockLost,
|
PutPartBeforeLockLost,
|
||||||
PutPartAfterRename,
|
PutPartAfterRename,
|
||||||
@@ -159,72 +156,6 @@ impl Drop for MultipartCommitBarrier {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
struct NewMultipartUploadCommitObservationState {
|
|
||||||
bucket: String,
|
|
||||||
object: String,
|
|
||||||
committed: AtomicBool,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) struct NewMultipartUploadCommitObservation {
|
|
||||||
state: Arc<NewMultipartUploadCommitObservationState>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
static NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION: std::sync::OnceLock<
|
|
||||||
std::sync::Mutex<Option<Arc<NewMultipartUploadCommitObservationState>>>,
|
|
||||||
> = std::sync::OnceLock::new();
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl NewMultipartUploadCommitObservation {
|
|
||||||
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
|
||||||
let state = Arc::new(NewMultipartUploadCommitObservationState {
|
|
||||||
bucket: bucket.to_string(),
|
|
||||||
object: object.to_string(),
|
|
||||||
committed: AtomicBool::new(false),
|
|
||||||
});
|
|
||||||
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("new multipart upload commit observation mutex should not poison");
|
|
||||||
assert!(slot.is_none(), "new multipart upload commit observation must be unique");
|
|
||||||
*slot = Some(Arc::clone(&state));
|
|
||||||
Self { state }
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn committed(&self) -> bool {
|
|
||||||
self.state.committed.load(Ordering::Acquire)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl Drop for NewMultipartUploadCommitObservation {
|
|
||||||
fn drop(&mut self) {
|
|
||||||
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("new multipart upload commit observation mutex should not poison");
|
|
||||||
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
|
||||||
*slot = None;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
fn observe_new_multipart_upload_commit(bucket: &str, object: &str) {
|
|
||||||
let state = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("new multipart upload commit observation mutex should not poison")
|
|
||||||
.as_ref()
|
|
||||||
.filter(|state| state.bucket == bucket && state.object == object)
|
|
||||||
.cloned();
|
|
||||||
if let Some(state) = state {
|
|
||||||
state.committed.store(true, Ordering::Release);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
|
async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
|
||||||
let barrier = {
|
let barrier = {
|
||||||
@@ -1684,30 +1615,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
|
|
||||||
let upload_path = Self::get_multipart_upload_dir(bucket, object, upload_uuid.as_str(), opts.data_movement);
|
let upload_path = Self::get_multipart_upload_dir(bucket, object, upload_uuid.as_str(), opts.data_movement);
|
||||||
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
|
||||||
pause_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost).await;
|
|
||||||
if _object_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
|
|
||||||
return Err(StorageError::NamespaceLockQuorumUnavailable {
|
|
||||||
mode: "new_multipart_upload_commit",
|
|
||||||
bucket: bucket.to_string(),
|
|
||||||
object: object.to_string(),
|
|
||||||
required: 1,
|
|
||||||
achieved: 0,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
if opts
|
|
||||||
.namespace_lock_fence
|
|
||||||
.as_ref()
|
|
||||||
.is_some_and(NamespaceLockFence::is_lock_lost)
|
|
||||||
{
|
|
||||||
return Err(StorageError::NamespaceLockQuorumUnavailable {
|
|
||||||
mode: "new_multipart_upload_outer_lock",
|
|
||||||
bucket: bucket.to_string(),
|
|
||||||
object: object.to_string(),
|
|
||||||
required: 1,
|
|
||||||
achieved: 0,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
|
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
|
||||||
Self::write_unique_file_info(
|
Self::write_unique_file_info(
|
||||||
&shuffle_disks,
|
&shuffle_disks,
|
||||||
@@ -1719,8 +1626,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
|
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
|
||||||
#[cfg(test)]
|
|
||||||
observe_new_multipart_upload_commit(bucket, object);
|
|
||||||
|
|
||||||
// evalDisks
|
// evalDisks
|
||||||
|
|
||||||
|
|||||||
@@ -2483,7 +2483,6 @@ impl SetDisks {
|
|||||||
})
|
})
|
||||||
.await?,
|
.await?,
|
||||||
);
|
);
|
||||||
notify_put_object_commit_namespace_acquired(bucket, object);
|
|
||||||
}
|
}
|
||||||
#[cfg(not(any(test, feature = "test-util")))]
|
#[cfg(not(any(test, feature = "test-util")))]
|
||||||
{
|
{
|
||||||
@@ -4631,7 +4630,6 @@ struct PutObjectCommitBarrierState {
|
|||||||
arrived: tokio::sync::Notify,
|
arrived: tokio::sync::Notify,
|
||||||
release: tokio::sync::Notify,
|
release: tokio::sync::Notify,
|
||||||
namespace_pending: tokio::sync::Notify,
|
namespace_pending: tokio::sync::Notify,
|
||||||
namespace_acquired: std::sync::atomic::AtomicBool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
@@ -4653,7 +4651,6 @@ impl PutObjectCommitBarrier {
|
|||||||
arrived: tokio::sync::Notify::new(),
|
arrived: tokio::sync::Notify::new(),
|
||||||
release: tokio::sync::Notify::new(),
|
release: tokio::sync::Notify::new(),
|
||||||
namespace_pending: tokio::sync::Notify::new(),
|
namespace_pending: tokio::sync::Notify::new(),
|
||||||
namespace_acquired: std::sync::atomic::AtomicBool::new(false),
|
|
||||||
});
|
});
|
||||||
let mut slot = PUT_OBJECT_COMMIT_BARRIER
|
let mut slot = PUT_OBJECT_COMMIT_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||||
@@ -4688,10 +4685,6 @@ impl PutObjectCommitBarrier {
|
|||||||
.await
|
.await
|
||||||
.expect("put object should wait for the namespace lock after leaving the commit barrier");
|
.expect("put object should wait for the namespace lock after leaving the commit barrier");
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn namespace_acquired(&self) -> bool {
|
|
||||||
self.state.namespace_acquired.load(std::sync::atomic::Ordering::Acquire)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
@@ -4748,22 +4741,6 @@ fn notify_put_object_commit_namespace_pending(bucket: &str, object: &str) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
|
||||||
fn notify_put_object_commit_namespace_acquired(bucket: &str, object: &str) {
|
|
||||||
let barrier = PUT_OBJECT_COMMIT_BARRIER
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
|
||||||
.lock()
|
|
||||||
.expect("put object commit barrier mutex should not poison")
|
|
||||||
.iter()
|
|
||||||
.find(|barrier| {
|
|
||||||
barrier.bucket == bucket && barrier.object == object && barrier.pause == PutObjectCommitPause::BeforeNamespace
|
|
||||||
})
|
|
||||||
.cloned();
|
|
||||||
if let Some(barrier) = barrier {
|
|
||||||
barrier.namespace_acquired.store(true, std::sync::atomic::Ordering::Release);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
struct DeleteObjectCommitBarrierState {
|
struct DeleteObjectCommitBarrierState {
|
||||||
bucket: String,
|
bucket: String,
|
||||||
@@ -4773,7 +4750,7 @@ struct DeleteObjectCommitBarrierState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pub(crate) struct DeleteObjectCommitBarrier {
|
struct DeleteObjectCommitBarrier {
|
||||||
state: Arc<DeleteObjectCommitBarrierState>,
|
state: Arc<DeleteObjectCommitBarrierState>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -4783,7 +4760,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
impl DeleteObjectCommitBarrier {
|
impl DeleteObjectCommitBarrier {
|
||||||
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
fn install(bucket: &str, object: &str) -> Self {
|
||||||
let state = Arc::new(DeleteObjectCommitBarrierState {
|
let state = Arc::new(DeleteObjectCommitBarrierState {
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
object: object.to_string(),
|
object: object.to_string(),
|
||||||
@@ -4799,13 +4776,13 @@ impl DeleteObjectCommitBarrier {
|
|||||||
Self { state }
|
Self { state }
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn wait_until_paused(&self) {
|
async fn wait_until_paused(&self) {
|
||||||
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
|
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
|
||||||
.await
|
.await
|
||||||
.expect("delete object should reach the deterministic commit barrier");
|
.expect("delete object should reach the deterministic commit barrier");
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn release(&self) {
|
fn release(&self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -5900,7 +5877,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
if dobj.version_id.is_none() && (version_suspended || versioned) {
|
if dobj.version_id.is_none() && (version_suspended || versioned) {
|
||||||
vr.mod_time = Some(OffsetDateTime::now_utc());
|
vr.mod_time = Some(OffsetDateTime::now_utc());
|
||||||
vr.deleted = true;
|
vr.deleted = true;
|
||||||
vr.mark_deleted = true;
|
|
||||||
if versioned {
|
if versioned {
|
||||||
vr.version_id = Some(Uuid::new_v4());
|
vr.version_id = Some(Uuid::new_v4());
|
||||||
}
|
}
|
||||||
@@ -11660,7 +11636,6 @@ mod transition_upload_integrity_tests {
|
|||||||
crate::data_movement::SourceCleanupBucketFence {
|
crate::data_movement::SourceCleanupBucketFence {
|
||||||
expected_incarnation_id: None,
|
expected_incarnation_id: None,
|
||||||
lifecycle_guard: Some(&bucket_guard),
|
lifecycle_guard: Some(&bucket_guard),
|
||||||
..Default::default()
|
|
||||||
},
|
},
|
||||||
"test_data_movement",
|
"test_data_movement",
|
||||||
)
|
)
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
@@ -151,7 +151,7 @@ pub(crate) mod init_format;
|
|||||||
pub(crate) mod list_objects;
|
pub(crate) mod list_objects;
|
||||||
mod multipart;
|
mod multipart;
|
||||||
mod object;
|
mod object;
|
||||||
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence};
|
pub(crate) use object::ObjectLockDiagGuard;
|
||||||
pub use object::{
|
pub use object::{
|
||||||
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
|
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
|
||||||
SnapshotConsistencyError,
|
SnapshotConsistencyError,
|
||||||
|
|||||||
@@ -400,7 +400,7 @@ impl ECStore {
|
|||||||
object: &str,
|
object: &str,
|
||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
) -> Result<MultipartUploadResult> {
|
) -> Result<MultipartUploadResult> {
|
||||||
self.handle_new_multipart_upload_with_pool_idx(bucket, object, opts, None)
|
self.handle_new_multipart_upload_with_pool_idx(bucket, object, opts)
|
||||||
.await
|
.await
|
||||||
.map(|(res, _, _)| res)
|
.map(|(res, _, _)| res)
|
||||||
}
|
}
|
||||||
@@ -410,22 +410,20 @@ impl ECStore {
|
|||||||
bucket: &str,
|
bucket: &str,
|
||||||
object: &str,
|
object: &str,
|
||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
|
||||||
) -> Result<(MultipartUploadResult, usize, Option<Uuid>)> {
|
) -> Result<(MultipartUploadResult, usize, Option<Uuid>)> {
|
||||||
check_new_multipart_args(bucket, object)?;
|
check_new_multipart_args(bucket, object)?;
|
||||||
let (mut opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?;
|
let (opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?;
|
||||||
|
let opts = &opts;
|
||||||
|
|
||||||
if self.single_pool() {
|
if self.single_pool() {
|
||||||
self.apply_decommission_target_mutation_fence(0, object, &mut opts, mutation_fence)
|
|
||||||
.await;
|
|
||||||
return self.pools[0]
|
return self.pools[0]
|
||||||
.new_multipart_upload(bucket, object, &opts)
|
.new_multipart_upload(bucket, object, opts)
|
||||||
.await
|
.await
|
||||||
.map(|res| (res, 0, opts.expected_bucket_incarnation_id));
|
.map(|res| (res, 0, opts.expected_bucket_incarnation_id));
|
||||||
}
|
}
|
||||||
|
|
||||||
if opts.data_movement && opts.version_id.is_some() {
|
if opts.data_movement && opts.version_id.is_some() {
|
||||||
let idx = self.select_data_movement_pool_idx(bucket, object, -1, &opts, false).await?;
|
let idx = self.select_data_movement_pool_idx(bucket, object, -1, opts, false).await?;
|
||||||
if idx == opts.src_pool_idx {
|
if idx == opts.src_pool_idx {
|
||||||
return Err(StorageError::DataMovementOverwriteErr(
|
return Err(StorageError::DataMovementOverwriteErr(
|
||||||
bucket.to_owned(),
|
bucket.to_owned(),
|
||||||
@@ -433,9 +431,7 @@ impl ECStore {
|
|||||||
opts.version_id.clone().unwrap_or_default(),
|
opts.version_id.clone().unwrap_or_default(),
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence)
|
let res = self.pools[idx].new_multipart_upload(bucket, object, opts).await?;
|
||||||
.await;
|
|
||||||
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
|
||||||
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -458,9 +454,7 @@ impl ECStore {
|
|||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
if !res.uploads.is_empty() {
|
if !res.uploads.is_empty() {
|
||||||
self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence)
|
let res = self.pools[idx].new_multipart_upload(bucket, object, opts).await?;
|
||||||
.await;
|
|
||||||
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
|
||||||
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -473,9 +467,7 @@ impl ECStore {
|
|||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence)
|
let res = self.pools[idx].new_multipart_upload(bucket, object, opts).await?;
|
||||||
.await;
|
|
||||||
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
|
||||||
Ok((res, idx, opts.expected_bucket_incarnation_id))
|
Ok((res, idx, opts.expected_bucket_incarnation_id))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -712,14 +704,13 @@ impl ECStore {
|
|||||||
|
|
||||||
pub(crate) async fn complete_multipart_upload_for_data_movement(
|
pub(crate) async fn complete_multipart_upload_for_data_movement(
|
||||||
self: Arc<Self>,
|
self: Arc<Self>,
|
||||||
target: (usize, Option<&ObjectLockDiagGuard>),
|
target_pool_idx: usize,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
object: &str,
|
object: &str,
|
||||||
upload_id: &str,
|
upload_id: &str,
|
||||||
uploaded_parts: Vec<CompletePart>,
|
uploaded_parts: Vec<CompletePart>,
|
||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
) -> Result<ObjectInfo> {
|
) -> Result<ObjectInfo> {
|
||||||
let (target_pool_idx, mutation_fence) = target;
|
|
||||||
check_complete_multipart_args(bucket, object, upload_id)?;
|
check_complete_multipart_args(bucket, object, upload_id)?;
|
||||||
if !opts.data_movement {
|
if !opts.data_movement {
|
||||||
return Err(Error::other("targeted multipart completion requires data_movement options"));
|
return Err(Error::other("targeted multipart completion requires data_movement options"));
|
||||||
@@ -748,8 +739,6 @@ impl ECStore {
|
|||||||
snapshot.add_lock_fences(&mut opts);
|
snapshot.add_lock_fences(&mut opts);
|
||||||
opts.object_lock_config_snapshot = Some(snapshot);
|
opts.object_lock_config_snapshot = Some(snapshot);
|
||||||
}
|
}
|
||||||
self.apply_decommission_target_mutation_fence(target_pool_idx, object, &mut opts, mutation_fence)
|
|
||||||
.await;
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pause_data_movement_multipart_before_selected_completion(bucket).await;
|
pause_data_movement_multipart_before_selected_completion(bucket).await;
|
||||||
let pool = self
|
let pool = self
|
||||||
|
|||||||
@@ -32,13 +32,12 @@ use crate::bucket::metadata_sys::{
|
|||||||
use crate::bucket::object_lock::objectlock_sys::{
|
use crate::bucket::object_lock::objectlock_sys::{
|
||||||
check_object_lock_for_deletion_with_state, ensure_recursive_force_delete_allowed_for_state,
|
check_object_lock_for_deletion_with_state, ensure_recursive_force_delete_allowed_for_state,
|
||||||
};
|
};
|
||||||
use crate::bucket::replication::{DeleteReplicationConfigSnapshot, ReplicationObjectBridge};
|
use crate::bucket::replication::ReplicationObjectBridge;
|
||||||
use crate::bucket::versioning::VersioningApi;
|
|
||||||
use crate::disk::OldCurrentSize;
|
use crate::disk::OldCurrentSize;
|
||||||
use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot};
|
use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot};
|
||||||
use crate::set_disk::{
|
use crate::set_disk::{
|
||||||
SetDisks, get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
|
get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
|
||||||
is_lock_optimization_enabled, is_object_lock_diag_enabled, same_distributed_lock_domain,
|
is_lock_optimization_enabled, is_object_lock_diag_enabled,
|
||||||
};
|
};
|
||||||
use crate::storage_api_contracts::{
|
use crate::storage_api_contracts::{
|
||||||
namespace::NamespaceLocking as _,
|
namespace::NamespaceLocking as _,
|
||||||
@@ -353,8 +352,6 @@ impl fmt::Display for ObjectLockDiagMode {
|
|||||||
|
|
||||||
pub(crate) struct ObjectLockDiagGuard {
|
pub(crate) struct ObjectLockDiagGuard {
|
||||||
guard: rustfs_lock::NamespaceLockGuard,
|
guard: rustfs_lock::NamespaceLockGuard,
|
||||||
#[cfg(test)]
|
|
||||||
test_namespace_lock_fence: Option<NamespaceLockFence>,
|
|
||||||
enabled: bool,
|
enabled: bool,
|
||||||
op: &'static str,
|
op: &'static str,
|
||||||
bucket: Option<String>,
|
bucket: Option<String>,
|
||||||
@@ -376,8 +373,6 @@ impl ObjectLockDiagGuard {
|
|||||||
) -> Self {
|
) -> Self {
|
||||||
Self {
|
Self {
|
||||||
guard,
|
guard,
|
||||||
#[cfg(test)]
|
|
||||||
test_namespace_lock_fence: None,
|
|
||||||
enabled,
|
enabled,
|
||||||
op,
|
op,
|
||||||
bucket,
|
bucket,
|
||||||
@@ -398,115 +393,6 @@ impl ObjectLockDiagGuard {
|
|||||||
pub(crate) fn is_lock_lost(&self) -> bool {
|
pub(crate) fn is_lock_lost(&self) -> bool {
|
||||||
self.guard.is_lock_lost()
|
self.guard.is_lock_lost()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
|
||||||
opts.ensure_namespace_lock_fence();
|
|
||||||
if let Some(signal) = self.lock_lost_signal() {
|
|
||||||
opts.add_namespace_lock_lost_signal(signal);
|
|
||||||
}
|
|
||||||
#[cfg(test)]
|
|
||||||
if let Some(fence) = self.test_namespace_lock_fence.as_ref() {
|
|
||||||
opts.add_namespace_lock_fence_for_test(fence);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
|
||||||
pub(crate) enum DecommissionMutationFenceTestPhase {
|
|
||||||
Migration,
|
|
||||||
SourceCleanup,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
struct DecommissionMutationFenceLossState {
|
|
||||||
bucket: String,
|
|
||||||
object: String,
|
|
||||||
phase: DecommissionMutationFenceTestPhase,
|
|
||||||
fence: NamespaceLockFence,
|
|
||||||
loss_handle: Arc<std::sync::atomic::AtomicBool>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) struct DecommissionMutationFenceLossHook {
|
|
||||||
state: Arc<DecommissionMutationFenceLossState>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
static DECOMMISSION_MUTATION_FENCE_LOSS_HOOK: std::sync::OnceLock<
|
|
||||||
std::sync::Mutex<Option<Arc<DecommissionMutationFenceLossState>>>,
|
|
||||||
> = std::sync::OnceLock::new();
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl DecommissionMutationFenceLossHook {
|
|
||||||
pub(crate) fn install(bucket: &str, object: &str, phase: DecommissionMutationFenceTestPhase) -> Self {
|
|
||||||
let (fence, loss_handle) = NamespaceLockFence::loss_handle_for_test();
|
|
||||||
let state = Arc::new(DecommissionMutationFenceLossState {
|
|
||||||
bucket: bucket.to_string(),
|
|
||||||
object: object.to_string(),
|
|
||||||
phase,
|
|
||||||
fence,
|
|
||||||
loss_handle,
|
|
||||||
});
|
|
||||||
let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("decommission mutation fence loss hooks should not poison");
|
|
||||||
assert!(slot.is_none(), "decommission mutation fence loss hook must be unique");
|
|
||||||
*slot = Some(Arc::clone(&state));
|
|
||||||
Self { state }
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn mark_lost(&self) {
|
|
||||||
self.state.loss_handle.store(true, Ordering::Release);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl Drop for DecommissionMutationFenceLossHook {
|
|
||||||
fn drop(&mut self) {
|
|
||||||
let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("decommission mutation fence loss hooks should not poison");
|
|
||||||
if slot.as_ref().is_some_and(|hook| Arc::ptr_eq(hook, &self.state)) {
|
|
||||||
*slot = None;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
fn decommission_mutation_fence_for_test(
|
|
||||||
bucket: &str,
|
|
||||||
object: &str,
|
|
||||||
phase: DecommissionMutationFenceTestPhase,
|
|
||||||
) -> Option<NamespaceLockFence> {
|
|
||||||
DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("decommission mutation fence loss hooks should not poison")
|
|
||||||
.as_ref()
|
|
||||||
.filter(|hook| hook.bucket == bucket && hook.object == object && hook.phase == phase)
|
|
||||||
.map(|hook| hook.fence.clone())
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) struct SourceCleanupMutationFence {
|
|
||||||
guard: ObjectLockDiagGuard,
|
|
||||||
source_lock_covered: bool,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl SourceCleanupMutationFence {
|
|
||||||
pub(crate) fn source_lock_covered(&self) -> bool {
|
|
||||||
self.source_lock_covered
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn is_lock_lost(&self) -> bool {
|
|
||||||
self.guard.is_lock_lost()
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
|
||||||
self.guard.add_namespace_lock_fence(opts);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Opaque write-lock guard for the RestoreObject accept path; see
|
/// Opaque write-lock guard for the RestoreObject accept path; see
|
||||||
@@ -524,7 +410,10 @@ impl RestoreAcceptGuard {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
pub fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
|
||||||
self.0.add_namespace_lock_fence(opts);
|
opts.ensure_namespace_lock_fence();
|
||||||
|
if let Some(signal) = self.0.lock_lost_signal() {
|
||||||
|
opts.add_namespace_lock_lost_signal(signal);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -801,6 +690,16 @@ impl SelectObjectSnapshotLockLossWake {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// LockRegistry clones its canonical client Arc for each endpoint host, so an
|
||||||
|
// exact Arc set identifies one distributed namespace-lock quorum domain.
|
||||||
|
fn same_distributed_lock_domain(left: &[Arc<dyn rustfs_lock::LockClient>], right: &[Arc<dyn rustfs_lock::LockClient>]) -> bool {
|
||||||
|
left.iter()
|
||||||
|
.all(|left_client| right.iter().any(|right_client| Arc::ptr_eq(left_client, right_client)))
|
||||||
|
&& right
|
||||||
|
.iter()
|
||||||
|
.all(|right_client| left.iter().any(|left_client| Arc::ptr_eq(left_client, right_client)))
|
||||||
|
}
|
||||||
|
|
||||||
impl AsyncRead for SelectObjectSnapshotReader {
|
impl AsyncRead for SelectObjectSnapshotReader {
|
||||||
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||||
if self.lock_loss_wake.poll_lost(cx) || self.lease.is_lost() {
|
if self.lock_loss_wake.poll_lost(cx) || self.lease.is_lost() {
|
||||||
@@ -906,7 +805,7 @@ fn resolve_latest_object_access(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool {
|
fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool {
|
||||||
(opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement
|
opts.versioned && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -914,8 +813,6 @@ struct DeleteAfterObjectLockSnapshotBarrierState {
|
|||||||
bucket: String,
|
bucket: String,
|
||||||
arrived: tokio::sync::Notify,
|
arrived: tokio::sync::Notify,
|
||||||
release: tokio::sync::Notify,
|
release: tokio::sync::Notify,
|
||||||
namespace_pending: tokio::sync::Notify,
|
|
||||||
namespace_acquired: AtomicBool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -935,8 +832,6 @@ impl DeleteAfterObjectLockSnapshotBarrier {
|
|||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
arrived: tokio::sync::Notify::new(),
|
arrived: tokio::sync::Notify::new(),
|
||||||
release: tokio::sync::Notify::new(),
|
release: tokio::sync::Notify::new(),
|
||||||
namespace_pending: tokio::sync::Notify::new(),
|
|
||||||
namespace_acquired: AtomicBool::new(false),
|
|
||||||
});
|
});
|
||||||
let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
|
let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
@@ -954,18 +849,6 @@ impl DeleteAfterObjectLockSnapshotBarrier {
|
|||||||
pub(crate) fn release(&self) {
|
pub(crate) fn release(&self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn release_and_wait_until_namespace_pending(&self) {
|
|
||||||
let namespace_pending = self.state.namespace_pending.notified();
|
|
||||||
self.release();
|
|
||||||
tokio::time::timeout(Duration::from_secs(5), namespace_pending)
|
|
||||||
.await
|
|
||||||
.expect("delete should proceed to its namespace lock after leaving the snapshot barrier");
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn namespace_acquired(&self) -> bool {
|
|
||||||
self.state.namespace_acquired.load(Ordering::Acquire)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -990,97 +873,6 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
|
|||||||
.as_ref()
|
.as_ref()
|
||||||
.filter(|state| state.bucket == bucket)
|
.filter(|state| state.bucket == bucket)
|
||||||
.cloned();
|
.cloned();
|
||||||
if let Some(state) = state {
|
|
||||||
state.arrived.notify_one();
|
|
||||||
state.release.notified().await;
|
|
||||||
state.namespace_pending.notify_one();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
fn notify_delete_namespace_acquired(bucket: &str) {
|
|
||||||
let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("delete snapshot barrier mutex should not poison")
|
|
||||||
.as_ref()
|
|
||||||
.filter(|state| state.bucket == bucket)
|
|
||||||
.cloned();
|
|
||||||
if let Some(state) = state {
|
|
||||||
state.namespace_acquired.store(true, Ordering::Release);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
struct VersionedDeleteMarkerCommitBarrierState {
|
|
||||||
bucket: String,
|
|
||||||
object: String,
|
|
||||||
arrived: tokio::sync::Notify,
|
|
||||||
release: tokio::sync::Notify,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) struct VersionedDeleteMarkerCommitBarrier {
|
|
||||||
state: Arc<VersionedDeleteMarkerCommitBarrierState>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
static VERSIONED_DELETE_MARKER_COMMIT_BARRIER: std::sync::OnceLock<
|
|
||||||
std::sync::Mutex<Option<Arc<VersionedDeleteMarkerCommitBarrierState>>>,
|
|
||||||
> = std::sync::OnceLock::new();
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl VersionedDeleteMarkerCommitBarrier {
|
|
||||||
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
|
||||||
let state = Arc::new(VersionedDeleteMarkerCommitBarrierState {
|
|
||||||
bucket: bucket.to_string(),
|
|
||||||
object: object.to_string(),
|
|
||||||
arrived: tokio::sync::Notify::new(),
|
|
||||||
release: tokio::sync::Notify::new(),
|
|
||||||
});
|
|
||||||
let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("versioned delete-marker commit barrier mutex should not poison");
|
|
||||||
assert!(slot.is_none(), "versioned delete-marker commit barrier must be unique");
|
|
||||||
*slot = Some(Arc::clone(&state));
|
|
||||||
Self { state }
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) async fn wait_until_paused(&self) {
|
|
||||||
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
|
|
||||||
.await
|
|
||||||
.expect("versioned DELETE should reach the post-marker-commit barrier");
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn release(&self) {
|
|
||||||
self.state.release.notify_one();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl Drop for VersionedDeleteMarkerCommitBarrier {
|
|
||||||
fn drop(&mut self) {
|
|
||||||
self.state.release.notify_one();
|
|
||||||
let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("versioned delete-marker commit barrier mutex should not poison");
|
|
||||||
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
|
||||||
*slot = None;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
async fn pause_versioned_delete_marker_after_commit(bucket: &str, object: &str) {
|
|
||||||
let state = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("versioned delete-marker commit barrier mutex should not poison")
|
|
||||||
.as_ref()
|
|
||||||
.filter(|state| state.bucket == bucket && state.object == object)
|
|
||||||
.cloned();
|
|
||||||
if let Some(state) = state {
|
if let Some(state) = state {
|
||||||
state.arrived.notify_one();
|
state.arrived.notify_one();
|
||||||
state.release.notified().await;
|
state.release.notified().await;
|
||||||
@@ -1121,160 +913,6 @@ fn writer_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions
|
|||||||
lookup_opts
|
lookup_opts
|
||||||
}
|
}
|
||||||
|
|
||||||
fn delete_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
|
|
||||||
let mut lookup_opts = writer_pool_lookup_opts(opts, no_lock);
|
|
||||||
lookup_opts.skip_decommissioned = opts.data_movement;
|
|
||||||
lookup_opts
|
|
||||||
}
|
|
||||||
|
|
||||||
fn should_delete_from_all_pools(opts: &ObjectOptions, pool_count: usize) -> bool {
|
|
||||||
pool_count > 0 && (!opts.versioned && !opts.version_suspended || opts.version_id.is_some())
|
|
||||||
}
|
|
||||||
|
|
||||||
fn batch_delete_creates_latest_marker(object: &ObjectToDelete, delete_config_snapshot: &DeleteReplicationConfigSnapshot) -> bool {
|
|
||||||
if object.version_id.is_some() {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
|
|
||||||
let object_name = decode_dir_object(&object.object_name);
|
|
||||||
let (versioned, version_suspended) = delete_config_snapshot.versioning_config().delete_state(&object_name);
|
|
||||||
versioned || version_suspended
|
|
||||||
}
|
|
||||||
|
|
||||||
fn batch_delete_targets_pool(creates_latest_marker: bool, marker_target_pool_idx: Option<usize>, pool_idx: usize) -> bool {
|
|
||||||
!creates_latest_marker || marker_target_pool_idx == Some(pool_idx)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
struct BatchDeletePoolErrorInjectionState {
|
|
||||||
bucket: String,
|
|
||||||
pool_idx: usize,
|
|
||||||
errors: std::collections::HashMap<String, Error>,
|
|
||||||
observed: std::sync::atomic::AtomicUsize,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
pub(crate) struct BatchDeletePoolErrorInjection {
|
|
||||||
state: Arc<BatchDeletePoolErrorInjectionState>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
static BATCH_DELETE_POOL_ERROR_INJECTION: std::sync::OnceLock<std::sync::Mutex<Option<Arc<BatchDeletePoolErrorInjectionState>>>> =
|
|
||||||
std::sync::OnceLock::new();
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl BatchDeletePoolErrorInjection {
|
|
||||||
pub(crate) fn install(bucket: &str, pool_idx: usize, errors: Vec<(String, Error)>) -> Self {
|
|
||||||
let state = Arc::new(BatchDeletePoolErrorInjectionState {
|
|
||||||
bucket: bucket.to_string(),
|
|
||||||
pool_idx,
|
|
||||||
errors: errors.into_iter().collect(),
|
|
||||||
observed: std::sync::atomic::AtomicUsize::new(0),
|
|
||||||
});
|
|
||||||
let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("batch delete pool error injection mutex should not poison");
|
|
||||||
assert!(slot.is_none(), "batch delete pool error injection must be unique");
|
|
||||||
*slot = Some(Arc::clone(&state));
|
|
||||||
Self { state }
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) fn observed(&self) -> usize {
|
|
||||||
self.state.observed.load(Ordering::Acquire)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl Drop for BatchDeletePoolErrorInjection {
|
|
||||||
fn drop(&mut self) {
|
|
||||||
let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("batch delete pool error injection mutex should not poison");
|
|
||||||
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
|
||||||
*slot = None;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
fn inject_batch_delete_pool_errors(
|
|
||||||
bucket: &str,
|
|
||||||
pool_idx: usize,
|
|
||||||
object_names: &[String],
|
|
||||||
result: &mut (Vec<DeletedObject>, Vec<Option<Error>>),
|
|
||||||
) {
|
|
||||||
let state = BATCH_DELETE_POOL_ERROR_INJECTION
|
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
|
||||||
.lock()
|
|
||||||
.expect("batch delete pool error injection mutex should not poison")
|
|
||||||
.as_ref()
|
|
||||||
.filter(|state| state.bucket == bucket && state.pool_idx == pool_idx)
|
|
||||||
.cloned();
|
|
||||||
let Some(state) = state else {
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
|
|
||||||
for (idx, object_name) in object_names.iter().enumerate() {
|
|
||||||
let Some(error) = state.errors.get(object_name) else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
if result.1[idx].is_none() && result.0[idx].found {
|
|
||||||
result.1[idx] = Some(error.clone());
|
|
||||||
state.observed.fetch_add(1, Ordering::AcqRel);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn resolve_batch_delete_pool_results<'a>(
|
|
||||||
initial_error: Option<Error>,
|
|
||||||
pool_results: impl IntoIterator<Item = (&'a DeletedObject, &'a Option<Error>)>,
|
|
||||||
) -> (Option<DeletedObject>, Option<Error>, bool) {
|
|
||||||
let mut failure = initial_error.map(|err| (None, err));
|
|
||||||
let mut deleted = None;
|
|
||||||
let mut fallback: Option<(DeletedObject, Option<Error>)> = None;
|
|
||||||
let mut attempted = false;
|
|
||||||
|
|
||||||
for (pool_delete, pool_error) in pool_results {
|
|
||||||
attempted = true;
|
|
||||||
match pool_error {
|
|
||||||
Some(err) if is_err_object_not_found(err) || is_err_version_not_found(err) => {
|
|
||||||
if fallback.as_ref().is_none_or(|(_, error)| error.is_none()) {
|
|
||||||
fallback = Some(((*pool_delete).clone(), Some(err.clone())));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Some(err) => {
|
|
||||||
if failure.is_none() {
|
|
||||||
failure = Some((Some((*pool_delete).clone()), err.clone()));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
None if pool_delete.found => {
|
|
||||||
if deleted.is_none() {
|
|
||||||
deleted = Some((*pool_delete).clone());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
None => {
|
|
||||||
if fallback.is_none() {
|
|
||||||
fallback = Some(((*pool_delete).clone(), None));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if let Some((failed_delete, err)) = failure {
|
|
||||||
return (failed_delete, Some(err), attempted);
|
|
||||||
}
|
|
||||||
if let Some(deleted) = deleted {
|
|
||||||
return (Some(deleted), None, attempted);
|
|
||||||
}
|
|
||||||
if let Some((deleted, err)) = fallback {
|
|
||||||
return (Some(deleted), err, attempted);
|
|
||||||
}
|
|
||||||
|
|
||||||
(None, None, attempted)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
|
fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
|
||||||
let mut lookup_opts = opts.clone();
|
let mut lookup_opts = opts.clone();
|
||||||
lookup_opts.skip_decommissioned = true;
|
lookup_opts.skip_decommissioned = true;
|
||||||
@@ -1895,89 +1533,6 @@ impl ECStore {
|
|||||||
)))
|
)))
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn acquire_decommission_object_mutation_fence(
|
|
||||||
&self,
|
|
||||||
bucket: &str,
|
|
||||||
object: &str,
|
|
||||||
) -> Result<ObjectLockDiagGuard> {
|
|
||||||
if self.ctx.lock_manager().is_disabled() {
|
|
||||||
return Err(Error::other("decommission object migration requires namespace locking"));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
let test_namespace_lock_fence =
|
|
||||||
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::Migration);
|
|
||||||
let object = encode_dir_object(object);
|
|
||||||
let mut opts = ObjectOptions::default();
|
|
||||||
let guard = self
|
|
||||||
.acquire_object_read_lock_if_needed("decommission_object", bucket, &object, &mut opts)
|
|
||||||
.await?
|
|
||||||
.ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence"))?;
|
|
||||||
#[cfg(test)]
|
|
||||||
let guard = {
|
|
||||||
let mut guard = guard;
|
|
||||||
guard.test_namespace_lock_fence = test_namespace_lock_fence;
|
|
||||||
guard
|
|
||||||
};
|
|
||||||
Ok(guard)
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(super) async fn apply_decommission_target_mutation_fence(
|
|
||||||
&self,
|
|
||||||
target_pool_idx: usize,
|
|
||||||
object: &str,
|
|
||||||
opts: &mut ObjectOptions,
|
|
||||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
|
||||||
) {
|
|
||||||
let Some(mutation_fence) = mutation_fence else {
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
|
|
||||||
mutation_fence.add_namespace_lock_fence(opts);
|
|
||||||
let fixed_set = self.pools.first().and_then(|pool| pool.disk_set.first());
|
|
||||||
let target_set = self.pools.get(target_pool_idx).map(|pool| pool.get_disks_by_key(object));
|
|
||||||
opts.no_lock = match (fixed_set, target_set) {
|
|
||||||
(Some(fixed), Some(target)) => fixed.shares_namespace_lock_domain(&target).await,
|
|
||||||
_ => false,
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) async fn acquire_decommission_source_cleanup_fence(
|
|
||||||
&self,
|
|
||||||
bucket: &str,
|
|
||||||
object: &str,
|
|
||||||
source_set: &SetDisks,
|
|
||||||
) -> Result<SourceCleanupMutationFence> {
|
|
||||||
if self.ctx.lock_manager().is_disabled() {
|
|
||||||
return Err(Error::other("decommission source cleanup requires namespace locking"));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
crate::data_movement::notify_source_cleanup_mutation_fence_pending(bucket, object);
|
|
||||||
#[cfg(test)]
|
|
||||||
let test_namespace_lock_fence =
|
|
||||||
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup);
|
|
||||||
let object = encode_dir_object(object);
|
|
||||||
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
|
|
||||||
let source_lock_covered = fixed_set.shares_namespace_lock_domain(source_set).await;
|
|
||||||
// Lock order: fixed store mutation domain first; source cleanup takes its
|
|
||||||
// hashed source-domain lock second only when this guard does not cover it.
|
|
||||||
let guard = self
|
|
||||||
.acquire_object_write_lock("decommission_source_cleanup", bucket, &object)
|
|
||||||
.await?;
|
|
||||||
#[cfg(test)]
|
|
||||||
let guard = {
|
|
||||||
let mut guard = guard;
|
|
||||||
guard.test_namespace_lock_fence = test_namespace_lock_fence;
|
|
||||||
guard
|
|
||||||
};
|
|
||||||
|
|
||||||
Ok(SourceCleanupMutationFence {
|
|
||||||
guard,
|
|
||||||
source_lock_covered,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
pub(crate) async fn acquire_all_object_read_locks(
|
pub(crate) async fn acquire_all_object_read_locks(
|
||||||
&self,
|
&self,
|
||||||
op: &'static str,
|
op: &'static str,
|
||||||
@@ -2431,17 +1986,14 @@ impl ECStore {
|
|||||||
object: &str,
|
object: &str,
|
||||||
data: &mut PutObjReader,
|
data: &mut PutObjReader,
|
||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
|
||||||
) -> Result<(usize, Result<ObjectInfo>)> {
|
) -> Result<(usize, Result<ObjectInfo>)> {
|
||||||
if !opts.data_movement {
|
if !opts.data_movement {
|
||||||
return Err(Error::other("data movement PUT requires data_movement options"));
|
return Err(Error::other("data movement PUT requires data_movement options"));
|
||||||
}
|
}
|
||||||
let (object, mut opts) = self.prepare_put_object(bucket, object, opts).await?;
|
let (object, opts) = self.prepare_put_object(bucket, object, opts).await?;
|
||||||
let idx = self
|
let idx = self
|
||||||
.select_put_object_pool_idx(bucket, object.as_str(), data.size(), &opts)
|
.select_put_object_pool_idx(bucket, object.as_str(), data.size(), &opts)
|
||||||
.await?;
|
.await?;
|
||||||
self.apply_decommission_target_mutation_fence(idx, object.as_str(), &mut opts, mutation_fence)
|
|
||||||
.await;
|
|
||||||
let result = self.pools[idx]
|
let result = self.pools[idx]
|
||||||
.put_object_with_old_current_size(bucket, &object, data, &opts)
|
.put_object_with_old_current_size(bucket, &object, data, &opts)
|
||||||
.await
|
.await
|
||||||
@@ -2894,10 +2446,6 @@ impl ECStore {
|
|||||||
} else {
|
} else {
|
||||||
None
|
None
|
||||||
};
|
};
|
||||||
#[cfg(test)]
|
|
||||||
if _object_lock_guard.is_some() {
|
|
||||||
notify_delete_namespace_acquired(bucket);
|
|
||||||
}
|
|
||||||
if let Some(trigger) = opts.lifecycle_delete_all.as_ref() {
|
if let Some(trigger) = opts.lifecycle_delete_all.as_ref() {
|
||||||
let configs = delete_all_configs.as_ref().ok_or(StorageError::PreconditionFailed)?;
|
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)?;
|
let expected_bucket_incarnation_id = opts.expected_bucket_incarnation_id.ok_or(StorageError::PreconditionFailed)?;
|
||||||
@@ -2931,7 +2479,7 @@ impl ECStore {
|
|||||||
return Ok(ObjectInfo::default());
|
return Ok(ObjectInfo::default());
|
||||||
}
|
}
|
||||||
|
|
||||||
let gopts = delete_pool_lookup_opts(&opts, true);
|
let gopts = writer_pool_lookup_opts(&opts, true);
|
||||||
|
|
||||||
if opts.data_movement {
|
if opts.data_movement {
|
||||||
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
|
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
|
||||||
@@ -3036,8 +2584,6 @@ impl ECStore {
|
|||||||
Err(err) if is_err_object_not_found(&err) && should_create_delete_marker_for_missing_object(&opts) => {
|
Err(err) if is_err_object_not_found(&err) && should_create_delete_marker_for_missing_object(&opts) => {
|
||||||
let target_pool_idx = self.get_pool_idx_no_lock(bucket, object, 0).await?;
|
let target_pool_idx = self.get_pool_idx_no_lock(bucket, object, 0).await?;
|
||||||
let mut obj = self.pools[target_pool_idx].delete_object(bucket, object, opts).await?;
|
let mut obj = self.pools[target_pool_idx].delete_object(bucket, object, opts).await?;
|
||||||
#[cfg(test)]
|
|
||||||
pause_versioned_delete_marker_after_commit(bucket, object).await;
|
|
||||||
obj.name = decode_dir_object(object);
|
obj.name = decode_dir_object(object);
|
||||||
return Ok(obj);
|
return Ok(obj);
|
||||||
}
|
}
|
||||||
@@ -3076,7 +2622,7 @@ impl ECStore {
|
|||||||
None
|
None
|
||||||
};
|
};
|
||||||
|
|
||||||
if should_delete_from_all_pools(&opts, errs.len()) {
|
if !errs.is_empty() && !opts.versioned && !opts.version_suspended {
|
||||||
let mut obj = match self.delete_object_from_all_pools(bucket, object, &opts, errs).await {
|
let mut obj = match self.delete_object_from_all_pools(bucket, object, &opts, errs).await {
|
||||||
Ok(obj) => obj,
|
Ok(obj) => obj,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -3100,8 +2646,6 @@ impl ECStore {
|
|||||||
|
|
||||||
match pool.delete_object(bucket, object, opts.clone()).await {
|
match pool.delete_object(bucket, object, opts.clone()).await {
|
||||||
Ok(res) => {
|
Ok(res) => {
|
||||||
#[cfg(test)]
|
|
||||||
pause_versioned_delete_marker_after_commit(bucket, object).await;
|
|
||||||
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
|
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
|
||||||
commit_prepared_tier_delete_journal_entry(api, je).await;
|
commit_prepared_tier_delete_journal_entry(api, je).await;
|
||||||
}
|
}
|
||||||
@@ -3232,104 +2776,30 @@ impl ECStore {
|
|||||||
Ok(guards) => guards,
|
Ok(guards) => guards,
|
||||||
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
|
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
|
||||||
};
|
};
|
||||||
#[cfg(test)]
|
|
||||||
if !_object_lock_guards.is_empty() {
|
|
||||||
notify_delete_namespace_acquired(bucket);
|
|
||||||
}
|
|
||||||
|
|
||||||
let delete_config_snapshot = opts
|
|
||||||
.delete_replication_config_snapshot
|
|
||||||
.as_deref()
|
|
||||||
.expect("batch delete replication config snapshot should be loaded");
|
|
||||||
let latest_marker_objects = objects
|
|
||||||
.iter()
|
|
||||||
.map(|object| batch_delete_creates_latest_marker(object, delete_config_snapshot))
|
|
||||||
.collect::<Vec<_>>();
|
|
||||||
let marker_target_results = join_all(objects.iter().zip(&latest_marker_objects).map(
|
|
||||||
|(object, creates_marker)| async move {
|
|
||||||
if *creates_marker {
|
|
||||||
Some(self.get_pool_idx_no_lock(bucket, &object.object_name, 0).await)
|
|
||||||
} else {
|
|
||||||
None
|
|
||||||
}
|
|
||||||
},
|
|
||||||
))
|
|
||||||
.await;
|
|
||||||
let mut marker_target_pool_indices = Vec::with_capacity(objects.len());
|
|
||||||
for (idx, target_result) in marker_target_results.into_iter().enumerate() {
|
|
||||||
match target_result {
|
|
||||||
Some(Ok(pool_idx)) => marker_target_pool_indices.push(Some(pool_idx)),
|
|
||||||
Some(Err(err)) => {
|
|
||||||
del_errs[idx] = Some(err);
|
|
||||||
marker_target_pool_indices.push(None);
|
|
||||||
}
|
|
||||||
None => marker_target_pool_indices.push(None),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
let mut futures = Vec::with_capacity(self.pools.len());
|
let mut futures = Vec::with_capacity(self.pools.len());
|
||||||
|
|
||||||
for pool in self.pools.iter() {
|
for pool in self.pools.iter() {
|
||||||
if self.is_pool_rebalancing(pool.pool_idx).await {
|
if self.is_pool_rebalancing(pool.pool_idx).await {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone()));
|
||||||
let (object_indices, pool_objects): (Vec<_>, Vec<_>) = objects
|
|
||||||
.iter()
|
|
||||||
.enumerate()
|
|
||||||
.filter(|(idx, _)| {
|
|
||||||
batch_delete_targets_pool(latest_marker_objects[*idx], marker_target_pool_indices[*idx], pool.pool_idx)
|
|
||||||
})
|
|
||||||
.map(|(idx, object)| (idx, object.clone()))
|
|
||||||
.unzip();
|
|
||||||
if pool_objects.is_empty() {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
let pool_opts = opts.clone();
|
|
||||||
futures.push(async move {
|
|
||||||
#[cfg(test)]
|
|
||||||
let pool_object_names = pool_objects
|
|
||||||
.iter()
|
|
||||||
.map(|object| object.object_name.clone())
|
|
||||||
.collect::<Vec<_>>();
|
|
||||||
let result = pool.delete_objects(bucket, pool_objects, pool_opts).await;
|
|
||||||
#[cfg(test)]
|
|
||||||
let result = {
|
|
||||||
let mut result = result;
|
|
||||||
inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, &mut result);
|
|
||||||
result
|
|
||||||
};
|
|
||||||
(object_indices, result)
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let results = join_all(futures).await;
|
let results = join_all(futures).await;
|
||||||
|
|
||||||
for idx in 0..del_objects.len() {
|
for idx in 0..del_objects.len() {
|
||||||
let pool_results = results.iter().filter_map(|(object_indices, (dels, errs))| {
|
for (dels, errs) in results.iter() {
|
||||||
let pool_object_idx = object_indices.binary_search(&idx).ok()?;
|
if errs[idx].is_none() && dels[idx].found {
|
||||||
Some((&dels[pool_object_idx], &errs[pool_object_idx]))
|
del_errs[idx] = None;
|
||||||
});
|
del_objects[idx] = dels[idx].clone();
|
||||||
let (deleted, error, attempted) = resolve_batch_delete_pool_results(del_errs[idx].take(), pool_results);
|
break;
|
||||||
if let Some(deleted) = deleted {
|
}
|
||||||
del_objects[idx] = deleted;
|
|
||||||
}
|
|
||||||
del_errs[idx] = error;
|
|
||||||
|
|
||||||
if !attempted && del_errs[idx].is_none() && latest_marker_objects[idx] {
|
if del_errs[idx].is_none() {
|
||||||
del_objects[idx] = DeletedObject {
|
del_errs[idx] = errs[idx].clone();
|
||||||
object_name: objects[idx].object_name.clone(),
|
del_objects[idx] = dels[idx].clone();
|
||||||
version_id: objects[idx].version_id,
|
}
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
del_errs[idx] = Some(StorageError::ObjectNotFound(bucket.to_owned(), objects[idx].object_name.clone()));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
for (idx, object) in objects.iter().enumerate() {
|
|
||||||
if del_errs[idx].is_none() && del_objects[idx].delete_marker {
|
|
||||||
pause_versioned_delete_marker_after_commit(bucket, &object.object_name).await;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3904,80 +3374,6 @@ mod tests {
|
|||||||
assert!(!same_distributed_lock_domain(&[first, second], &[other]));
|
assert!(!same_distributed_lock_domain(&[first, second], &[other]));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn decommission_fence_covers_dist_sets_with_same_clients_despite_different_namespaces() {
|
|
||||||
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
|
|
||||||
let (_dirs, original_sets) = make_local_two_set_sets_with_ctx(Arc::clone(&ctx)).await;
|
|
||||||
let mut second_set = (*original_sets.disk_set[1]).clone();
|
|
||||||
second_set.lockers = original_sets.disk_set[0].lockers.clone();
|
|
||||||
let mut sets = (*original_sets).clone();
|
|
||||||
sets.disk_set[1] = Arc::new(second_set);
|
|
||||||
let sets = Arc::new(sets);
|
|
||||||
ctx.update_erasure_type(SetupType::DistErasure).await;
|
|
||||||
|
|
||||||
assert!(
|
|
||||||
sets.disk_set[0]
|
|
||||||
.lockers
|
|
||||||
.iter()
|
|
||||||
.zip(&sets.disk_set[1].lockers)
|
|
||||||
.all(|(fixed, hashed)| Arc::ptr_eq(fixed, hashed)),
|
|
||||||
"the regression requires identical distributed lock clients"
|
|
||||||
);
|
|
||||||
assert_ne!(sets.disk_set[0].set_index, sets.disk_set[1].set_index);
|
|
||||||
|
|
||||||
let pool_config = sets.endpoints.clone();
|
|
||||||
let store = new_prepared_reader_test_store_from_pools(vec![Arc::clone(&sets)], vec![pool_config], ctx);
|
|
||||||
let object = (0..1_000)
|
|
||||||
.map(|index| format!("decommission-dist-domain-{index}.bin"))
|
|
||||||
.find(|candidate| Arc::ptr_eq(&sets.get_disks_by_key(candidate), &sets.disk_set[1]))
|
|
||||||
.expect("a key should hash to the second set namespace");
|
|
||||||
let mutation_fence = store
|
|
||||||
.acquire_decommission_object_mutation_fence("bucket", &object)
|
|
||||||
.await
|
|
||||||
.expect("the fixed distributed mutation fence should be acquired");
|
|
||||||
let target_lock = sets.disk_set[1]
|
|
||||||
.new_ns_lock("bucket", &object)
|
|
||||||
.await
|
|
||||||
.expect("the hashed-set namespace lock should be created");
|
|
||||||
let target_err = target_lock
|
|
||||||
.get_write_lock(Duration::from_millis(50))
|
|
||||||
.await
|
|
||||||
.expect_err("the fixed read fence must conflict through the shared clients");
|
|
||||||
assert!(matches!(target_err, rustfs_lock::LockError::Timeout { .. }));
|
|
||||||
|
|
||||||
let mut put_opts = ObjectOptions::default();
|
|
||||||
store
|
|
||||||
.apply_decommission_target_mutation_fence(0, &object, &mut put_opts, Some(&mutation_fence))
|
|
||||||
.await;
|
|
||||||
assert!(put_opts.no_lock, "migration target PUT must reuse the covering fixed fence");
|
|
||||||
|
|
||||||
let mut multipart_opts = ObjectOptions::default();
|
|
||||||
store
|
|
||||||
.apply_decommission_target_mutation_fence(0, &object, &mut multipart_opts, Some(&mutation_fence))
|
|
||||||
.await;
|
|
||||||
assert!(multipart_opts.no_lock, "migration target multipart must reuse the covering fixed fence");
|
|
||||||
drop(mutation_fence);
|
|
||||||
|
|
||||||
let cleanup_object = (0..1_000)
|
|
||||||
.map(|index| format!("decommission-dist-cleanup-{index}.bin"))
|
|
||||||
.find(|candidate| Arc::ptr_eq(&sets.get_disks_by_key(candidate), &sets.disk_set[1]))
|
|
||||||
.expect("a cleanup key should hash to the second set namespace");
|
|
||||||
let source_fence = store
|
|
||||||
.acquire_decommission_source_cleanup_fence("bucket", &cleanup_object, sets.disk_set[1].as_ref())
|
|
||||||
.await
|
|
||||||
.expect("the fixed distributed cleanup fence should be acquired");
|
|
||||||
assert!(source_fence.source_lock_covered(), "source cleanup must reuse the covering fixed fence");
|
|
||||||
let source_lock = sets.disk_set[1]
|
|
||||||
.new_ns_lock("bucket", &cleanup_object)
|
|
||||||
.await
|
|
||||||
.expect("the source-set namespace lock should be created");
|
|
||||||
let source_err = source_lock
|
|
||||||
.get_read_lock(Duration::from_millis(50))
|
|
||||||
.await
|
|
||||||
.expect_err("the fixed write fence must conflict through the shared clients");
|
|
||||||
assert!(matches!(source_err, rustfs_lock::LockError::Timeout { .. }));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn select_snapshot_version_matching_normalizes_null_and_uuid_forms() {
|
fn select_snapshot_version_matching_normalizes_null_and_uuid_forms() {
|
||||||
let nil = Uuid::nil();
|
let nil = Uuid::nil();
|
||||||
@@ -5037,159 +4433,6 @@ mod tests {
|
|||||||
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
|
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn ordinary_delete_lookup_includes_decommission_source_and_skips_rebalance_source() {
|
|
||||||
let lookup_opts = delete_pool_lookup_opts(&ObjectOptions::default(), true);
|
|
||||||
|
|
||||||
assert!(lookup_opts.no_lock);
|
|
||||||
assert!(!lookup_opts.skip_decommissioned);
|
|
||||||
assert!(lookup_opts.skip_rebalancing);
|
|
||||||
|
|
||||||
let explicit_version = delete_pool_lookup_opts(
|
|
||||||
&ObjectOptions {
|
|
||||||
versioned: true,
|
|
||||||
version_id: Some(uuid::Uuid::new_v4().to_string()),
|
|
||||||
..Default::default()
|
|
||||||
},
|
|
||||||
true,
|
|
||||||
);
|
|
||||||
assert!(!explicit_version.skip_decommissioned);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn delete_fans_out_for_unversioned_and_explicit_version_mutations() {
|
|
||||||
assert!(should_delete_from_all_pools(&ObjectOptions::default(), 1));
|
|
||||||
assert!(should_delete_from_all_pools(
|
|
||||||
&ObjectOptions {
|
|
||||||
versioned: true,
|
|
||||||
version_id: Some(uuid::Uuid::new_v4().to_string()),
|
|
||||||
..Default::default()
|
|
||||||
},
|
|
||||||
2,
|
|
||||||
));
|
|
||||||
assert!(!should_delete_from_all_pools(
|
|
||||||
&ObjectOptions {
|
|
||||||
versioned: true,
|
|
||||||
..Default::default()
|
|
||||||
},
|
|
||||||
1,
|
|
||||||
));
|
|
||||||
assert!(!should_delete_from_all_pools(&ObjectOptions::default(), 0));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn batch_delete_identifies_only_latest_versioned_markers() {
|
|
||||||
let versioned = DeleteReplicationConfigSnapshot::from_configs_for_test(
|
|
||||||
s3s::dto::VersioningConfiguration {
|
|
||||||
status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::ENABLED)),
|
|
||||||
..Default::default()
|
|
||||||
},
|
|
||||||
None,
|
|
||||||
);
|
|
||||||
let latest = ObjectToDelete {
|
|
||||||
object_name: "latest".to_string(),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
assert!(batch_delete_creates_latest_marker(&latest, &versioned));
|
|
||||||
assert!(!batch_delete_targets_pool(true, Some(1), 0));
|
|
||||||
assert!(batch_delete_targets_pool(true, Some(1), 1));
|
|
||||||
assert!(!batch_delete_targets_pool(true, Some(1), 2));
|
|
||||||
|
|
||||||
let explicit = ObjectToDelete {
|
|
||||||
object_name: "explicit".to_string(),
|
|
||||||
version_id: Some(uuid::Uuid::new_v4()),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
assert!(!batch_delete_creates_latest_marker(&explicit, &versioned));
|
|
||||||
assert!(batch_delete_targets_pool(false, Some(1), 0));
|
|
||||||
|
|
||||||
let unversioned = DeleteReplicationConfigSnapshot::default();
|
|
||||||
assert!(!batch_delete_creates_latest_marker(&latest, &unversioned));
|
|
||||||
assert!(batch_delete_targets_pool(false, None, 0));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn batch_delete_pool_failures_override_success_in_any_pool_order() {
|
|
||||||
let success = DeletedObject {
|
|
||||||
object_name: "object".to_string(),
|
|
||||||
found: true,
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
let source_errors = [
|
|
||||||
StorageError::ErasureWriteQuorum,
|
|
||||||
StorageError::NamespaceLockQuorumUnavailable {
|
|
||||||
mode: "delete_objects_commit",
|
|
||||||
bucket: "bucket".to_string(),
|
|
||||||
object: "object".to_string(),
|
|
||||||
required: 1,
|
|
||||||
achieved: 0,
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
for source_error in source_errors {
|
|
||||||
for source_first in [true, false] {
|
|
||||||
let failed = (DeletedObject::default(), Some(source_error.clone()));
|
|
||||||
let succeeded = (success.clone(), None);
|
|
||||||
let pool_results = if source_first {
|
|
||||||
vec![failed, succeeded]
|
|
||||||
} else {
|
|
||||||
vec![succeeded, failed]
|
|
||||||
};
|
|
||||||
|
|
||||||
let (_, error, attempted) =
|
|
||||||
resolve_batch_delete_pool_results(None, pool_results.iter().map(|(deleted, error)| (deleted, error)));
|
|
||||||
|
|
||||||
assert!(attempted);
|
|
||||||
assert_eq!(error, Some(source_error.clone()));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn batch_delete_ignores_missing_pool_only_after_another_pool_succeeds() {
|
|
||||||
let success = DeletedObject {
|
|
||||||
object_name: "object".to_string(),
|
|
||||||
found: true,
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
let missing_errors = [
|
|
||||||
StorageError::ObjectNotFound("bucket".to_string(), "object".to_string()),
|
|
||||||
StorageError::VersionNotFound("bucket".to_string(), "object".to_string(), "version".to_string()),
|
|
||||||
];
|
|
||||||
|
|
||||||
for missing_error in missing_errors {
|
|
||||||
let missing = (DeletedObject::default(), Some(missing_error.clone()));
|
|
||||||
for missing_first in [true, false] {
|
|
||||||
let succeeded = (success.clone(), None);
|
|
||||||
let pool_results = if missing_first {
|
|
||||||
vec![missing.clone(), succeeded]
|
|
||||||
} else {
|
|
||||||
vec![succeeded, missing.clone()]
|
|
||||||
};
|
|
||||||
let (deleted, error, attempted) =
|
|
||||||
resolve_batch_delete_pool_results(None, pool_results.iter().map(|(deleted, error)| (deleted, error)));
|
|
||||||
|
|
||||||
assert!(attempted);
|
|
||||||
let deleted = deleted.expect("successful pool result should be retained");
|
|
||||||
assert!(deleted.found);
|
|
||||||
assert_eq!(deleted.object_name, success.object_name.as_str());
|
|
||||||
assert!(error.is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
let missing_only = [missing];
|
|
||||||
let (_, error, attempted) =
|
|
||||||
resolve_batch_delete_pool_results(None, missing_only.iter().map(|(deleted, error)| (deleted, error)));
|
|
||||||
assert!(attempted);
|
|
||||||
assert_eq!(error, Some(missing_error));
|
|
||||||
}
|
|
||||||
|
|
||||||
let silent_missing = [(DeletedObject::default(), None)];
|
|
||||||
let (_, error, attempted) =
|
|
||||||
resolve_batch_delete_pool_results(None, silent_missing.iter().map(|(deleted, error)| (deleted, error)));
|
|
||||||
assert!(attempted);
|
|
||||||
assert!(error.is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() {
|
fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() {
|
||||||
let lookup_opts = data_movement_pool_lookup_opts(
|
let lookup_opts = data_movement_pool_lookup_opts(
|
||||||
|
|||||||
@@ -206,6 +206,13 @@ def check_runner_selection(root: Path) -> list[str]:
|
|||||||
return errors
|
return errors
|
||||||
|
|
||||||
|
|
||||||
|
def check_s3_tests_runner(root: Path) -> list[str]:
|
||||||
|
runner = (root / "scripts/s3-tests/run.sh").read_text()
|
||||||
|
if "--showlocals" in runner:
|
||||||
|
return ["scripts/s3-tests/run.sh: pytest failure diagnostics must not dump local values"]
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
def profile_selection(root: Path, profile: str) -> str:
|
def profile_selection(root: Path, profile: str) -> str:
|
||||||
if not re.fullmatch(r"e2e-[a-z0-9-]+", profile):
|
if not re.fullmatch(r"e2e-[a-z0-9-]+", profile):
|
||||||
raise ValueError(f"invalid e2e profile name: {profile}")
|
raise ValueError(f"invalid e2e profile name: {profile}")
|
||||||
@@ -272,6 +279,7 @@ def validate(root: Path) -> list[str]:
|
|||||||
errors.extend(check_e2e_modules(root))
|
errors.extend(check_e2e_modules(root))
|
||||||
errors.extend(check_fuzz_targets(root))
|
errors.extend(check_fuzz_targets(root))
|
||||||
errors.extend(check_runner_selection(root))
|
errors.extend(check_runner_selection(root))
|
||||||
|
errors.extend(check_s3_tests_runner(root))
|
||||||
errors.extend(check_profile_definitions(root))
|
errors.extend(check_profile_definitions(root))
|
||||||
return errors
|
return errors
|
||||||
|
|
||||||
@@ -341,6 +349,23 @@ class SelfTests(unittest.TestCase):
|
|||||||
)
|
)
|
||||||
self.assertEqual(len(check_fuzz_targets(root)), 1)
|
self.assertEqual(len(check_fuzz_targets(root)), 1)
|
||||||
|
|
||||||
|
def test_s3_runner_rejects_unbounded_failure_locals(self) -> None:
|
||||||
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
|
root = Path(tmp)
|
||||||
|
runner = root / "scripts/s3-tests/run.sh"
|
||||||
|
runner.parent.mkdir(parents=True)
|
||||||
|
runner.write_text("tox -- -vv -ra --tb=long\n")
|
||||||
|
self.assertEqual(check_s3_tests_runner(root), [])
|
||||||
|
runner.write_text("tox -- -vv -ra --showlocals --tb=long\n")
|
||||||
|
self.assertEqual(len(check_s3_tests_runner(root)), 1)
|
||||||
|
with (
|
||||||
|
mock.patch(__name__ + ".check_e2e_modules", return_value=[]),
|
||||||
|
mock.patch(__name__ + ".check_fuzz_targets", return_value=[]),
|
||||||
|
mock.patch(__name__ + ".check_runner_selection", return_value=[]),
|
||||||
|
mock.patch(__name__ + ".check_profile_definitions", return_value=[]),
|
||||||
|
):
|
||||||
|
self.assertEqual(len(validate(root)), 1)
|
||||||
|
|
||||||
def test_profile_listing_enforces_selection(self) -> None:
|
def test_profile_listing_enforces_selection(self) -> None:
|
||||||
with tempfile.TemporaryDirectory() as tmp:
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
root = Path(tmp)
|
root = Path(tmp)
|
||||||
@@ -411,7 +436,7 @@ def main() -> int:
|
|||||||
for error in errors:
|
for error in errors:
|
||||||
print(f"ERROR: {error}", file=sys.stderr)
|
print(f"ERROR: {error}", file=sys.stderr)
|
||||||
return 1
|
return 1
|
||||||
print("OK: e2e modules, runner selection, fuzz matrices, and profile guards are wired")
|
print("OK: e2e modules, runner selection, fuzz matrices, profiles, and bounded diagnostics are wired")
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1028,10 +1028,11 @@ else
|
|||||||
fi
|
fi
|
||||||
|
|
||||||
# Run tests from s3tests/functional
|
# Run tests from s3tests/functional
|
||||||
|
# Failure locals can contain multi-MiB request bodies; keep tracebacks without expanding local values.
|
||||||
set +e
|
set +e
|
||||||
S3TEST_CONF="${CONF_OUTPUT_PATH}" \
|
S3TEST_CONF="${CONF_OUTPUT_PATH}" \
|
||||||
tox -- \
|
tox -- \
|
||||||
-vv -ra --showlocals --tb=long \
|
-vv -ra --tb=long \
|
||||||
--maxfail="${MAXFAIL}" \
|
--maxfail="${MAXFAIL}" \
|
||||||
--timeout="${TEST_TIMEOUT}" \
|
--timeout="${TEST_TIMEOUT}" \
|
||||||
--junitxml="${ARTIFACTS_DIR}/junit.xml" \
|
--junitxml="${ARTIFACTS_DIR}/junit.xml" \
|
||||||
|
|||||||
Reference in New Issue
Block a user