mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 12:35:54 +00:00
fix(scanner): fence backup and cleanup mutations
Carry storage-owned publication scopes through backup CAS writes and observed cleanup deletes so cancellation and absolute lease deadlines retain the same safety boundary. Co-Authored-By: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -270,6 +270,22 @@ const OLD_DATA_CLEANUP_RECEIPT_FILE: &str = ".rustfs-old-data-cleanup-receipt.js
|
|||||||
const SCANNER_PUBLICATION_LEASE_FENCE_MAX_BYTES: usize = 64 * 1024;
|
const SCANNER_PUBLICATION_LEASE_FENCE_MAX_BYTES: usize = 64 * 1024;
|
||||||
const SCANNER_PUBLICATION_LEASE_FENCE_MAX_ENTRIES: usize = 256;
|
const SCANNER_PUBLICATION_LEASE_FENCE_MAX_ENTRIES: usize = 256;
|
||||||
|
|
||||||
|
fn begin_scanner_publication_delete_mutation(scope: Option<&crate::object_api::ScannerPublicationCommitScope>) -> Result<()> {
|
||||||
|
let Some(scope) = scope else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
if scope.state() == crate::object_api::ScannerPublicationCommitState::Admitted {
|
||||||
|
scope
|
||||||
|
.try_begin()
|
||||||
|
.map_err(|err| Error::other(format!("scanner publication delete scope cannot start: {err:?}")))?;
|
||||||
|
}
|
||||||
|
if !scope.can_commit() {
|
||||||
|
let _ = scope.mark_indeterminate();
|
||||||
|
return Err(StorageError::OperationCanceled);
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
fn take_scanner_publication_lease_tokens(user_defined: &mut HashMap<String, String>) -> Result<Option<HashMap<String, Uuid>>> {
|
fn take_scanner_publication_lease_tokens(user_defined: &mut HashMap<String, String>) -> Result<Option<HashMap<String, Uuid>>> {
|
||||||
let Some(encoded) = user_defined.remove(SCANNER_PUBLICATION_LEASE_FENCE_METADATA_KEY) else {
|
let Some(encoded) = user_defined.remove(SCANNER_PUBLICATION_LEASE_FENCE_METADATA_KEY) else {
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
@@ -7098,6 +7114,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
|
|
||||||
#[tracing::instrument(skip(self, opts))]
|
#[tracing::instrument(skip(self, opts))]
|
||||||
async fn delete_object(&self, bucket: &str, object: &str, mut opts: ObjectOptions) -> Result<ObjectInfo> {
|
async fn delete_object(&self, bucket: &str, object: &str, mut opts: ObjectOptions) -> Result<ObjectInfo> {
|
||||||
|
let _scope_outcome_guard = opts
|
||||||
|
.scanner_publication_commit_scope
|
||||||
|
.clone()
|
||||||
|
.map(ScannerPublicationCommitScopeGuard::new);
|
||||||
|
let scanner_publication_commit_scope = opts.scanner_publication_commit_scope.clone();
|
||||||
// Scanner cleanup carries the per-peer lease fence as transient
|
// Scanner cleanup carries the per-peer lease fence as transient
|
||||||
// request metadata. Consume it before any delete-prefix fanout so it
|
// request metadata. Consume it before any delete-prefix fanout so it
|
||||||
// cannot be persisted or treated as user metadata.
|
// cannot be persisted or treated as user metadata.
|
||||||
@@ -7192,6 +7213,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
}
|
}
|
||||||
delete_request.set_skip_tier_free_version();
|
delete_request.set_skip_tier_free_version();
|
||||||
}
|
}
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_object_version(bucket, object, &delete_request, false).await?;
|
self.delete_object_version(bucket, object, &delete_request, false).await?;
|
||||||
if let Some((_, deleted_object)) = replication_delete {
|
if let Some((_, deleted_object)) = replication_delete {
|
||||||
ReplicationLifecycleBridge::schedule_delete(bucket.to_string(), deleted_object).await;
|
ReplicationLifecycleBridge::schedule_delete(bucket.to_string(), deleted_object).await;
|
||||||
@@ -7206,6 +7228,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
delete_request.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
delete_request.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_object_version(bucket, object, &delete_request, false).await?;
|
self.delete_object_version(bucket, object, &delete_request, false).await?;
|
||||||
}
|
}
|
||||||
for version in &versions.free_versions {
|
for version in &versions.free_versions {
|
||||||
@@ -7217,10 +7240,14 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
delete_request.set_tier_free_version();
|
delete_request.set_tier_free_version();
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_object_version(bucket, object, &delete_request, false).await?;
|
self.delete_object_version(bucket, object, &delete_request, false).await?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if let Some(scope) = scanner_publication_commit_scope.as_ref() {
|
||||||
|
let _ = scope.mark_committed();
|
||||||
|
}
|
||||||
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||||
return Ok(ObjectInfo::default());
|
return Ok(ObjectInfo::default());
|
||||||
}
|
}
|
||||||
@@ -7228,10 +7255,14 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
self.validate_bucket_incarnation(bucket, expected_incarnation_id).await?;
|
self.validate_bucket_incarnation(bucket, expected_incarnation_id).await?;
|
||||||
}
|
}
|
||||||
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_prefix_with_scanner_publication_lease(bucket, object, scanner_publication_lease_tokens.as_ref())
|
self.delete_prefix_with_scanner_publication_lease(bucket, object, scanner_publication_lease_tokens.as_ref())
|
||||||
.await
|
.await
|
||||||
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
|
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
|
||||||
|
|
||||||
|
if let Some(scope) = scanner_publication_commit_scope.as_ref() {
|
||||||
|
let _ = scope.mark_committed();
|
||||||
|
}
|
||||||
self.invalidate_all_get_object_metadata_cache();
|
self.invalidate_all_get_object_metadata_cache();
|
||||||
return Ok(ObjectInfo::default());
|
return Ok(ObjectInfo::default());
|
||||||
}
|
}
|
||||||
@@ -7304,10 +7335,14 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_object_version(bucket, object, &dfi, false)
|
self.delete_object_version(bucket, object, &dfi, false)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
||||||
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||||
|
if let Some(scope) = scanner_publication_commit_scope.as_ref() {
|
||||||
|
let _ = scope.mark_committed();
|
||||||
|
}
|
||||||
return Ok(ObjectInfo::from_file_info(&dfi, bucket, object, opts.versioned || opts.version_suspended));
|
return Ok(ObjectInfo::from_file_info(&dfi, bucket, object, opts.versioned || opts.version_suspended));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -7381,6 +7416,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
};
|
};
|
||||||
|
|
||||||
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_object_version(bucket, object, &fi, should_force_delete_marker_for_missing_version(&opts))
|
self.delete_object_version(bucket, object, &fi, should_force_delete_marker_for_missing_version(&opts))
|
||||||
.await
|
.await
|
||||||
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
||||||
@@ -7392,6 +7428,9 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
oi.user_tags = Arc::clone(&goi.user_tags);
|
oi.user_tags = Arc::clone(&goi.user_tags);
|
||||||
oi.replication_decision = goi.replication_decision;
|
oi.replication_decision = goi.replication_decision;
|
||||||
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||||
|
if let Some(scope) = scanner_publication_commit_scope.as_ref() {
|
||||||
|
let _ = scope.mark_committed();
|
||||||
|
}
|
||||||
return Ok(oi);
|
return Ok(oi);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -7417,6 +7456,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
}
|
}
|
||||||
|
|
||||||
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
||||||
|
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
|
||||||
self.delete_object_version(bucket, object, &dfi, opts.delete_marker)
|
self.delete_object_version(bucket, object, &dfi, opts.delete_marker)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
||||||
@@ -7442,6 +7482,9 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
obj_info.delete_marker = true;
|
obj_info.delete_marker = true;
|
||||||
}
|
}
|
||||||
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||||
|
if let Some(scope) = scanner_publication_commit_scope.as_ref() {
|
||||||
|
let _ = scope.mark_committed();
|
||||||
|
}
|
||||||
Ok(obj_info)
|
Ok(obj_info)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -753,6 +753,32 @@ pub(crate) fn scanner_publication_epoch_changed(error: &EcstoreError) -> bool {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn delete_config_with_publication_scope_for_epoch<S>(
|
||||||
|
api: Arc<S>,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
mut opts: ScannerObjectOptions,
|
||||||
|
expected_epoch: u64,
|
||||||
|
scanner_publication_commit_scope: Option<ScannerPublicationCommitScope>,
|
||||||
|
) -> EcstoreResult<ScannerObjectInfo>
|
||||||
|
where
|
||||||
|
S: ScannerObjectIO + ScannerConfigObjectDelete,
|
||||||
|
{
|
||||||
|
let legacy_admission = if scanner_publication_commit_scope.is_none() {
|
||||||
|
Some(
|
||||||
|
scanner_publication_admission_for_epoch(api.clone(), expected_epoch)
|
||||||
|
.await
|
||||||
|
.ok_or_else(|| EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED))?,
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
};
|
||||||
|
opts.scanner_publication_commit_scope = scanner_publication_commit_scope;
|
||||||
|
let result = api.delete_config_object(bucket, object, opts).await;
|
||||||
|
drop(legacy_admission);
|
||||||
|
result
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn delete_config_with_publication_admission_for_epoch<S>(
|
pub(crate) async fn delete_config_with_publication_admission_for_epoch<S>(
|
||||||
api: Arc<S>,
|
api: Arc<S>,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -763,10 +789,7 @@ pub(crate) async fn delete_config_with_publication_admission_for_epoch<S>(
|
|||||||
where
|
where
|
||||||
S: ScannerObjectIO + ScannerConfigObjectDelete,
|
S: ScannerObjectIO + ScannerConfigObjectDelete,
|
||||||
{
|
{
|
||||||
let Some(_admission) = scanner_publication_admission_for_epoch(api.clone(), expected_epoch).await else {
|
delete_config_with_publication_scope_for_epoch(api, bucket, object, opts, expected_epoch, None).await
|
||||||
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
|
||||||
};
|
|
||||||
api.delete_config_object(bucket, object, opts).await
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Capture the storage-owned publication epoch without retaining the read
|
/// Capture the storage-owned publication epoch without retaining the read
|
||||||
|
|||||||
@@ -37,6 +37,7 @@ use crate::scanner_io::{
|
|||||||
scanner_maintenance_generation,
|
scanner_maintenance_generation,
|
||||||
};
|
};
|
||||||
use crate::sleeper::{SCANNER_SLEEPER, set_scanner_default_speed};
|
use crate::sleeper::{SCANNER_SLEEPER, set_scanner_default_speed};
|
||||||
|
use crate::storage_api::owner::ScannerPublicationCommitState;
|
||||||
use crate::{DataUsageInfo, ScannerActivityGuard, ScannerError, ScannerRuntimeGuard};
|
use crate::{DataUsageInfo, ScannerActivityGuard, ScannerError, ScannerRuntimeGuard};
|
||||||
use crate::{ScannerConfigObjectDelete, ScannerObjectIO, ScannerObjectOptions};
|
use crate::{ScannerConfigObjectDelete, ScannerObjectIO, ScannerObjectOptions};
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
@@ -455,7 +456,6 @@ fn data_usage_backup_due(data_usage_info: &DataUsageInfo) -> bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
#[allow(dead_code)]
|
|
||||||
async fn sync_data_usage_backup_from_primary(
|
async fn sync_data_usage_backup_from_primary(
|
||||||
ctx: &CancellationToken,
|
ctx: &CancellationToken,
|
||||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||||
@@ -463,7 +463,7 @@ async fn sync_data_usage_backup_from_primary(
|
|||||||
sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence(ctx, storeapi, None, None, None).await
|
sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence(ctx, storeapi, None, None, None).await
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(dead_code)]
|
#[cfg(test)]
|
||||||
async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence(
|
async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence(
|
||||||
ctx: &CancellationToken,
|
ctx: &CancellationToken,
|
||||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||||
@@ -471,25 +471,25 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence(
|
|||||||
remote_lease_deadline: Option<std::time::Instant>,
|
remote_lease_deadline: Option<std::time::Instant>,
|
||||||
scanner_publication_lease_fence: Option<&str>,
|
scanner_publication_lease_fence: Option<&str>,
|
||||||
) -> Result<(), EcstoreError> {
|
) -> Result<(), EcstoreError> {
|
||||||
sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_scope(
|
sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope(
|
||||||
ctx,
|
ctx,
|
||||||
storeapi,
|
storeapi,
|
||||||
expected_publication_epoch,
|
expected_publication_epoch,
|
||||||
remote_lease_deadline,
|
remote_lease_deadline,
|
||||||
scanner_publication_lease_fence,
|
scanner_publication_lease_fence,
|
||||||
Vec::new(),
|
&[],
|
||||||
Arc::new(AtomicBool::new(true)),
|
Arc::new(AtomicBool::new(true)),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_scope(
|
async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope(
|
||||||
ctx: &CancellationToken,
|
ctx: &CancellationToken,
|
||||||
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
|
||||||
expected_publication_epoch: Option<u64>,
|
expected_publication_epoch: Option<u64>,
|
||||||
remote_lease_deadline: Option<std::time::Instant>,
|
remote_lease_deadline: Option<std::time::Instant>,
|
||||||
scanner_publication_lease_fence: Option<&str>,
|
scanner_publication_lease_fence: Option<&str>,
|
||||||
remote_lease_tokens: Vec<Uuid>,
|
remote_lease_tokens: &[Uuid],
|
||||||
lease_release_safe: Arc<AtomicBool>,
|
lease_release_safe: Arc<AtomicBool>,
|
||||||
) -> Result<(), EcstoreError> {
|
) -> Result<(), EcstoreError> {
|
||||||
let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
|
let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
|
||||||
@@ -549,31 +549,32 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_s
|
|||||||
if remote_lease_deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) {
|
if remote_lease_deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) {
|
||||||
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
||||||
}
|
}
|
||||||
let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else {
|
let publication_scope = storeapi
|
||||||
if retry < SCANNER_PERSIST_CAS_RETRIES {
|
.scanner_data_usage_publication_commit_scope_with_release_flag(
|
||||||
continue;
|
read_epoch,
|
||||||
}
|
tokio::time::Instant::now()
|
||||||
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
.checked_add(data_usage_persist_timeout())
|
||||||
|
.unwrap_or_else(tokio::time::Instant::now)
|
||||||
|
.min(
|
||||||
|
remote_lease_deadline
|
||||||
|
.map(tokio::time::Instant::from_std)
|
||||||
|
.unwrap_or_else(|| tokio::time::Instant::now() + data_usage_persist_timeout()),
|
||||||
|
),
|
||||||
|
remote_lease_tokens.to_vec(),
|
||||||
|
Arc::clone(&lease_release_safe),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
let legacy_publication_admission = if publication_scope.is_none() {
|
||||||
|
let Some(admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else {
|
||||||
|
if retry < SCANNER_PERSIST_CAS_RETRIES {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
||||||
|
};
|
||||||
|
Some(admission)
|
||||||
|
} else {
|
||||||
|
None
|
||||||
};
|
};
|
||||||
let publication_scope = match expected_publication_epoch {
|
|
||||||
Some(expected_epoch) => {
|
|
||||||
storeapi
|
|
||||||
.scanner_data_usage_publication_commit_scope_with_release_flag(
|
|
||||||
expected_epoch,
|
|
||||||
usage_store::scanner_publication_scope_deadline(data_usage_persist_timeout(), remote_lease_deadline),
|
|
||||||
remote_lease_tokens.clone(),
|
|
||||||
Arc::clone(&lease_release_safe),
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
}
|
|
||||||
None => None,
|
|
||||||
};
|
|
||||||
if expected_publication_epoch.is_some() && publication_scope.is_none() {
|
|
||||||
if retry < SCANNER_PERSIST_CAS_RETRIES {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
|
|
||||||
}
|
|
||||||
let save_result = save_config_shared_with_preconditions_and_lease_fence_and_scope(
|
let save_result = save_config_shared_with_preconditions_and_lease_fence_and_scope(
|
||||||
storeapi.clone(),
|
storeapi.clone(),
|
||||||
&backup_path,
|
&backup_path,
|
||||||
@@ -584,19 +585,20 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_s
|
|||||||
publication_scope.clone(),
|
publication_scope.clone(),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
drop(legacy_publication_admission);
|
||||||
if let Some(scope) = publication_scope {
|
if let Some(scope) = publication_scope {
|
||||||
match scope.wait_for_completion().await {
|
match scope.wait_for_completion().await {
|
||||||
crate::storage_api::owner::ScannerPublicationCommitState::Committed
|
ScannerPublicationCommitState::Committed | ScannerPublicationCommitState::AbortedBeforeCommit => {}
|
||||||
| crate::storage_api::owner::ScannerPublicationCommitState::AbortedBeforeCommit => save_result,
|
ScannerPublicationCommitState::Indeterminate
|
||||||
crate::storage_api::owner::ScannerPublicationCommitState::Indeterminate
|
| ScannerPublicationCommitState::Admitted
|
||||||
| crate::storage_api::owner::ScannerPublicationCommitState::Admitted
|
| ScannerPublicationCommitState::InFlight => {
|
||||||
| crate::storage_api::owner::ScannerPublicationCommitState::InFlight => Err(EcstoreError::other(
|
return Err(EcstoreError::other(
|
||||||
"scanner backup publication commit scope did not reach a safe terminal state",
|
"scanner backup publication scope did not reach a safe terminal state",
|
||||||
)),
|
));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
} else {
|
|
||||||
save_result
|
|
||||||
}
|
}
|
||||||
|
save_result
|
||||||
};
|
};
|
||||||
|
|
||||||
match save_result {
|
match save_result {
|
||||||
@@ -1473,10 +1475,13 @@ async fn run_data_scanner_cycle_with_budget(
|
|||||||
mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await;
|
mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await;
|
||||||
return ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::DataMovement);
|
return ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::DataMovement);
|
||||||
};
|
};
|
||||||
let usage_persist_baseline_result = read_data_usage_persist_baseline(storeapi.clone()).await;
|
let usage_persist_baseline_result = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await;
|
||||||
drop(baseline_publication_guard);
|
drop(baseline_publication_guard);
|
||||||
let usage_persist_baseline = match usage_persist_baseline_result {
|
let usage_persist_baseline = match usage_persist_baseline_result {
|
||||||
Ok(baseline) => baseline,
|
Ok((data, revision)) => DataUsagePersistBaseline {
|
||||||
|
data: data.map(Bytes::from),
|
||||||
|
revision,
|
||||||
|
},
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
error!(
|
error!(
|
||||||
target: "rustfs::scanner",
|
target: "rustfs::scanner",
|
||||||
@@ -1516,20 +1521,6 @@ async fn run_data_scanner_cycle_with_budget(
|
|||||||
{
|
{
|
||||||
Some(ScannerCycleDeferReason::DataMovement)
|
Some(ScannerCycleDeferReason::DataMovement)
|
||||||
}
|
}
|
||||||
// A complete walk can still be retained as an observational snapshot
|
|
||||||
// when only the final activity proof was unavailable. It must not
|
|
||||||
// block the observation receiver: the authoritative publication
|
|
||||||
// fence remains enforced by the usage store and the cycle is advanced
|
|
||||||
// as partial without acknowledging dirty usage.
|
|
||||||
Ok(result)
|
|
||||||
if result.has_observational_snapshot()
|
|
||||||
&& matches!(
|
|
||||||
result.status,
|
|
||||||
ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable)
|
|
||||||
) =>
|
|
||||||
{
|
|
||||||
None
|
|
||||||
}
|
|
||||||
Ok(result) => final_data_usage_publication_defer_reason(storeapi.as_ref(), result.status).await,
|
Ok(result) => final_data_usage_publication_defer_reason(storeapi.as_ref(), result.status).await,
|
||||||
Err(_) => Some(ScannerCycleDeferReason::ActivityBaselineUnavailable),
|
Err(_) => Some(ScannerCycleDeferReason::ActivityBaselineUnavailable),
|
||||||
};
|
};
|
||||||
@@ -2879,7 +2870,12 @@ fn finalize_scanner_cycle_result(
|
|||||||
scan_cycle_result: crate::scanner_io::ScannerCycleResult,
|
scan_cycle_result: crate::scanner_io::ScannerCycleResult,
|
||||||
usage_persist_outcome: DataUsagePersistOutcome,
|
usage_persist_outcome: DataUsagePersistOutcome,
|
||||||
) -> (ScannerCycleOutcome, bool, Vec<ScannerDirtyUsageAcknowledgement>) {
|
) -> (ScannerCycleOutcome, bool, Vec<ScannerDirtyUsageAcknowledgement>) {
|
||||||
let completion_outcome = scanner_cycle_completion_outcome_for_result(&scan_cycle_result, usage_persist_outcome);
|
let completion_outcome = scanner_cycle_completion_outcome(
|
||||||
|
scan_cycle_result.status,
|
||||||
|
usage_persist_outcome,
|
||||||
|
scan_cycle_result.has_dirty_usage_to_acknowledge(),
|
||||||
|
scan_cycle_result.has_failed_dirty_usage(),
|
||||||
|
);
|
||||||
let pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work();
|
let pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work();
|
||||||
let durable_complete_snapshot = scan_cycle_result.status == ScannerCycleStatus::Complete
|
let durable_complete_snapshot = scan_cycle_result.status == ScannerCycleStatus::Complete
|
||||||
&& matches!(
|
&& matches!(
|
||||||
@@ -2894,34 +2890,6 @@ fn finalize_scanner_cycle_result(
|
|||||||
(completion_outcome, pending_maintenance_work, remote_dirty_usage_acknowledgements)
|
(completion_outcome, pending_maintenance_work, remote_dirty_usage_acknowledgements)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn scanner_cycle_completion_outcome_for_result(
|
|
||||||
scan_cycle_result: &crate::scanner_io::ScannerCycleResult,
|
|
||||||
usage_persist_outcome: DataUsagePersistOutcome,
|
|
||||||
) -> ScannerCycleOutcome {
|
|
||||||
let has_dirty_usage = scan_cycle_result.has_dirty_usage_to_acknowledge();
|
|
||||||
let has_failed_dirty_usage = scan_cycle_result.has_failed_dirty_usage();
|
|
||||||
if scan_cycle_result.has_observational_snapshot()
|
|
||||||
&& matches!(
|
|
||||||
scan_cycle_result.status,
|
|
||||||
ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable)
|
|
||||||
)
|
|
||||||
{
|
|
||||||
return match usage_persist_outcome {
|
|
||||||
DataUsagePersistOutcome::Saved
|
|
||||||
| DataUsagePersistOutcome::AlreadyDurable
|
|
||||||
| DataUsagePersistOutcome::PriorCycleDurable
|
|
||||||
| DataUsagePersistOutcome::Current
|
|
||||||
if !has_failed_dirty_usage =>
|
|
||||||
{
|
|
||||||
ScannerCycleOutcome::Partial
|
|
||||||
}
|
|
||||||
DataUsagePersistOutcome::Deferred(reason) => ScannerCycleOutcome::Deferred(reason),
|
|
||||||
_ => ScannerCycleOutcome::Failed,
|
|
||||||
};
|
|
||||||
}
|
|
||||||
scanner_cycle_completion_outcome(scan_cycle_result.status, usage_persist_outcome, has_dirty_usage, has_failed_dirty_usage)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Decide whether an incoming usage snapshot must be skipped as stale, given the local
|
/// Decide whether an incoming usage snapshot must be skipped as stale, given the local
|
||||||
/// wall clock `now`. Mirrors `stale_data_usage_persist_reason` in
|
/// wall clock `now`. Mirrors `stale_data_usage_persist_reason` in
|
||||||
/// `crates/ecstore/src/data_usage/mod.rs` — keep the two consistent.
|
/// `crates/ecstore/src/data_usage/mod.rs` — keep the two consistent.
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ fn remote_lease_expired(deadline: Option<std::time::Instant>) -> bool {
|
|||||||
deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline)
|
deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(super) fn scanner_publication_scope_deadline(
|
fn scanner_publication_scope_deadline(
|
||||||
persist_timeout: Duration,
|
persist_timeout: Duration,
|
||||||
remote_lease_deadline: Option<std::time::Instant>,
|
remote_lease_deadline: Option<std::time::Instant>,
|
||||||
) -> tokio::time::Instant {
|
) -> tokio::time::Instant {
|
||||||
@@ -53,84 +53,6 @@ pub(super) struct DataUsagePersistBaseline {
|
|||||||
pub(super) revision: DataUsageCacheRevision,
|
pub(super) revision: DataUsageCacheRevision,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Read the bytes used as the baseline for a usage publication while keeping
|
|
||||||
/// the v2 primary revision as the CAS fence. During an interrupted upgrade the
|
|
||||||
/// primary can be valid JSON without a baseline identity; in that case a
|
|
||||||
/// same-or-newer durable companion may still be used, but an older legacy
|
|
||||||
/// snapshot must not cross the primary's epoch fence.
|
|
||||||
pub(super) async fn read_data_usage_persist_baseline(
|
|
||||||
storeapi: Arc<impl ScannerObjectIO>,
|
|
||||||
) -> Result<DataUsagePersistBaseline, EcstoreError> {
|
|
||||||
let (primary, revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await?;
|
|
||||||
let Some(primary) = primary else {
|
|
||||||
for path in [
|
|
||||||
format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
|
||||||
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
|
|
||||||
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
|
||||||
] {
|
|
||||||
let (candidate, _) = read_config_with_revision(storeapi.clone(), &path).await?;
|
|
||||||
let Some(candidate) = candidate else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
let Ok(usage) = serde_json::from_slice::<DataUsageInfo>(&candidate) else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
if data_usage_info_has_persisted_baseline_identity(&usage) {
|
|
||||||
return Ok(DataUsagePersistBaseline {
|
|
||||||
data: Some(Bytes::from(candidate)),
|
|
||||||
revision,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return Ok(DataUsagePersistBaseline { data: None, revision });
|
|
||||||
};
|
|
||||||
|
|
||||||
let Ok(primary_info) = serde_json::from_slice::<DataUsageInfo>(&primary) else {
|
|
||||||
// Preserve the original bytes and revision. A completed scan may
|
|
||||||
// replace the invalid primary under this CAS fence; an observation
|
|
||||||
// will still reject it below because it has no verifiable identity.
|
|
||||||
return Ok(DataUsagePersistBaseline {
|
|
||||||
data: Some(Bytes::from(primary)),
|
|
||||||
revision,
|
|
||||||
});
|
|
||||||
};
|
|
||||||
if data_usage_info_has_persisted_baseline_identity(&primary_info) || data_usage_info_is_bootstrap_pending(&primary_info) {
|
|
||||||
return Ok(DataUsagePersistBaseline {
|
|
||||||
data: Some(Bytes::from(primary)),
|
|
||||||
revision,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
let invalid_primary_epoch = primary_info.scanner_epoch;
|
|
||||||
for path in [
|
|
||||||
format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
|
||||||
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
|
|
||||||
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
|
|
||||||
] {
|
|
||||||
let (candidate, _) = read_config_with_revision(storeapi.clone(), &path).await?;
|
|
||||||
let Some(candidate) = candidate else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
let Ok(usage) = serde_json::from_slice::<DataUsageInfo>(&candidate) else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
let candidate_epoch = usage.scanner_epoch.unwrap_or_default();
|
|
||||||
if data_usage_info_has_persisted_baseline_identity(&usage)
|
|
||||||
&& invalid_primary_epoch.is_none_or(|epoch| candidate_epoch >= epoch)
|
|
||||||
{
|
|
||||||
return Ok(DataUsagePersistBaseline {
|
|
||||||
data: Some(Bytes::from(candidate)),
|
|
||||||
revision,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(DataUsagePersistBaseline {
|
|
||||||
data: Some(Bytes::from(primary)),
|
|
||||||
revision,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Short-lived publication inputs captured for one usage persistence attempt.
|
/// Short-lived publication inputs captured for one usage persistence attempt.
|
||||||
/// Keeping the movement epoch, lease deadline, and target fence together makes
|
/// Keeping the movement epoch, lease deadline, and target fence together makes
|
||||||
/// it explicit that they are one proof rather than independent options.
|
/// it explicit that they are one proof rather than independent options.
|
||||||
@@ -388,8 +310,8 @@ where
|
|||||||
publication_epoch = Some(read_epoch);
|
publication_epoch = Some(read_epoch);
|
||||||
let authoritative_data = match next_baseline.as_ref() {
|
let authoritative_data = match next_baseline.as_ref() {
|
||||||
Some(baseline) => baseline.data.clone(),
|
Some(baseline) => baseline.data.clone(),
|
||||||
None => match read_data_usage_persist_baseline(storeapi.clone()).await {
|
None => match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await {
|
||||||
Ok(baseline) => baseline.data,
|
Ok((data, _)) => data.map(Bytes::from),
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
error!(
|
error!(
|
||||||
target: "rustfs::scanner",
|
target: "rustfs::scanner",
|
||||||
@@ -754,6 +676,8 @@ where
|
|||||||
expected_publication_epoch,
|
expected_publication_epoch,
|
||||||
remote_lease_deadline,
|
remote_lease_deadline,
|
||||||
scanner_publication_lease_fence.as_deref(),
|
scanner_publication_lease_fence.as_deref(),
|
||||||
|
&remote_lease_tokens,
|
||||||
|
Arc::clone(&lease_release_safe),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
if expected_publication_epoch.is_some() && !cleanup_ok {
|
if expected_publication_epoch.is_some() && !cleanup_ok {
|
||||||
@@ -777,6 +701,8 @@ where
|
|||||||
expected_publication_epoch,
|
expected_publication_epoch,
|
||||||
remote_lease_deadline,
|
remote_lease_deadline,
|
||||||
scanner_publication_lease_fence.as_deref(),
|
scanner_publication_lease_fence.as_deref(),
|
||||||
|
&remote_lease_tokens,
|
||||||
|
Arc::clone(&lease_release_safe),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
if expected_publication_epoch.is_some() && !cleanup_ok {
|
if expected_publication_epoch.is_some() && !cleanup_ok {
|
||||||
@@ -819,6 +745,8 @@ where
|
|||||||
expected_publication_epoch,
|
expected_publication_epoch,
|
||||||
remote_lease_deadline,
|
remote_lease_deadline,
|
||||||
scanner_publication_lease_fence.as_deref(),
|
scanner_publication_lease_fence.as_deref(),
|
||||||
|
&remote_lease_tokens,
|
||||||
|
Arc::clone(&lease_release_safe),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
if expected_publication_epoch.is_some() && !cleanup_ok {
|
if expected_publication_epoch.is_some() && !cleanup_ok {
|
||||||
@@ -836,13 +764,13 @@ where
|
|||||||
|
|
||||||
if backup_due {
|
if backup_due {
|
||||||
let done_save = Metrics::time(Metric::SaveUsage);
|
let done_save = Metrics::time(Metric::SaveUsage);
|
||||||
let backup_result = sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_scope(
|
let backup_result = sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope(
|
||||||
&ctx,
|
&ctx,
|
||||||
storeapi.clone(),
|
storeapi.clone(),
|
||||||
expected_publication_epoch,
|
expected_publication_epoch,
|
||||||
remote_lease_deadline,
|
remote_lease_deadline,
|
||||||
scanner_publication_lease_fence.as_deref(),
|
scanner_publication_lease_fence.as_deref(),
|
||||||
remote_lease_tokens.clone(),
|
&remote_lease_tokens,
|
||||||
Arc::clone(&lease_release_safe),
|
Arc::clone(&lease_release_safe),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
@@ -877,6 +805,8 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease(
|
|||||||
expected_publication_epoch: Option<u64>,
|
expected_publication_epoch: Option<u64>,
|
||||||
remote_lease_deadline: Option<std::time::Instant>,
|
remote_lease_deadline: Option<std::time::Instant>,
|
||||||
scanner_publication_lease_fence: Option<&str>,
|
scanner_publication_lease_fence: Option<&str>,
|
||||||
|
remote_lease_tokens: &[Uuid],
|
||||||
|
lease_release_safe: Arc<AtomicBool>,
|
||||||
) -> bool {
|
) -> bool {
|
||||||
if remote_lease_expired(remote_lease_deadline) {
|
if remote_lease_expired(remote_lease_deadline) {
|
||||||
return false;
|
return false;
|
||||||
@@ -945,7 +875,15 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease(
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
let result = delete_config_with_publication_admission_for_epoch(
|
let publication_scope = storeapi
|
||||||
|
.scanner_data_usage_publication_commit_scope_with_release_flag(
|
||||||
|
read_epoch,
|
||||||
|
scanner_publication_scope_deadline(data_usage_persist_timeout(), remote_lease_deadline),
|
||||||
|
remote_lease_tokens.to_vec(),
|
||||||
|
Arc::clone(&lease_release_safe),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
let result = crate::delete_config_with_publication_scope_for_epoch(
|
||||||
storeapi,
|
storeapi,
|
||||||
RUSTFS_META_BUCKET,
|
RUSTFS_META_BUCKET,
|
||||||
DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str(),
|
DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str(),
|
||||||
@@ -964,9 +902,23 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease(
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
read_epoch,
|
read_epoch,
|
||||||
|
publication_scope.clone(),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
|
let result = if let Some(scope) = publication_scope {
|
||||||
|
match scope.wait_for_completion().await {
|
||||||
|
ScannerPublicationCommitState::Committed | ScannerPublicationCommitState::AbortedBeforeCommit => result,
|
||||||
|
ScannerPublicationCommitState::Indeterminate
|
||||||
|
| ScannerPublicationCommitState::Admitted
|
||||||
|
| ScannerPublicationCommitState::InFlight => Err(EcstoreError::other(
|
||||||
|
"scanner publication cleanup scope did not reach a safe terminal state",
|
||||||
|
)),
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
result
|
||||||
|
};
|
||||||
|
|
||||||
match result {
|
match result {
|
||||||
Ok(_)
|
Ok(_)
|
||||||
| Err(
|
| Err(
|
||||||
|
|||||||
Reference in New Issue
Block a user