Compare commits

..

12 Commits

Author SHA1 Message Date
houseme 5ff5a7a06e Merge branch 'main' into cxymds/fix-1927-durable-checkpoint
Signed-off-by: houseme <housemecn@gmail.com>
2026-08-23 12:27:11 +08:00
overtrue c58accabe7 fix(ecstore): support Windows checkpoint CAS 2026-08-23 07:58:22 +08:00
马登山 4d92835d7b merge main and fix swift lint 2026-08-23 06:57:56 +08:00
overtrue 71db3327ad fix(heal): reset unverified checkpoint progress 2026-08-23 06:17:31 +08:00
overtrue 68e95d497e fix(heal): require current checkpoint digest 2026-08-23 04:19:04 +08:00
overtrue 8350d73be9 fix(storage): bound conditional file lock artifacts 2026-08-23 02:24:27 +08:00
overtrue 9a2a33a7eb fix(heal): canonicalize checkpoint integrity digest 2026-08-23 00:09:14 +08:00
cxymds 5df26a13ba Merge branch 'main' into cxymds/fix-1927-durable-checkpoint 2026-08-22 22:06:41 +08:00
马登山 e3a362989f fix(heal): atomically authenticate checkpoints 2026-08-22 22:05:23 +08:00
马登山 0cd2ae20e2 fix(heal): fail closed on tampered resume checkpoints 2026-08-22 15:32:12 +08:00
cxymds 55a7fa9f03 Merge branch 'main' into cxymds/fix-1927-durable-checkpoint 2026-08-22 11:25:52 +08:00
马登山 9f51f37a0d fix(heal): make resume checkpoints crash consistent 2026-08-21 22:04:23 +08:00
22 changed files with 944 additions and 1444 deletions
+6 -3
View File
@@ -39,10 +39,11 @@ jobs:
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
steps:
- name: Checkout repository
- name: Checkout main branch
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -88,10 +89,11 @@ jobs:
# either casing.
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout repository
- name: Checkout main branch
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -176,10 +178,11 @@ jobs:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout repository
- name: Checkout main branch
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
Generated
+1
View File
@@ -9564,6 +9564,7 @@ dependencies = [
"serde",
"serde_json",
"serial_test",
"sha2 0.11.0",
"temp-env",
"tempfile",
"thiserror 2.0.20",
+13 -263
View File
@@ -13,8 +13,6 @@
// limitations under the License.
use crate::bucket::replication::replication_state_from_filemeta;
#[cfg(test)]
use crate::bucket::utils::is_meta_bucketname;
use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::bucket::{
lifecycle::{
@@ -1349,53 +1347,6 @@ fn should_cleanup_decommission_source_entry(decommissioned: usize, total_version
decommissioned.saturating_add(expired) == total_versions
}
const DECOMMISSION_FREE_VERSION_MIGRATED_REASON: &str = "tier_free_version_migrated";
const DECOMMISSION_FREE_VERSION_CONSUMED_REASON: &str = "tier_free_version_already_consumed";
const DECOMMISSION_FREE_VERSION_RETAINED_REASON: &str = "tier_free_version_migration_failed";
const DECOMMISSION_FREE_VERSION_SWEEP_REASON: &str = "tier_free_version_unresolved_after_decommission";
const DECOMMISSION_FREE_VERSION_DISPOSITION_REASON: &str = "tier_free_version_disposition_recorded";
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
struct DecommissionFreeVersionDisposition {
migrated: usize,
consumed: usize,
retained: usize,
}
impl DecommissionFreeVersionDisposition {
fn record_migrated(&mut self) {
self.migrated += 1;
}
fn record_consumed(&mut self) {
self.consumed += 1;
}
fn record_retained(&mut self) {
self.retained += 1;
}
fn total(self) -> usize {
self.migrated.saturating_add(self.consumed).saturating_add(self.retained)
}
}
enum DecommissionFreeVersionAttempt {
Migrated,
Consumed,
CapacityFailure(Error),
Retry(Error),
}
fn classify_decommission_free_version_attempt(result: Result<()>) -> DecommissionFreeVersionAttempt {
match result {
Ok(()) => DecommissionFreeVersionAttempt::Migrated,
Err(err) if is_decommission_copy_cleanup_safe_error(&err) => DecommissionFreeVersionAttempt::Consumed,
Err(err) if is_decommission_target_capacity_error(&err) => DecommissionFreeVersionAttempt::CapacityFailure(err),
Err(err) => DecommissionFreeVersionAttempt::Retry(err),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[allow(
dead_code,
@@ -2664,12 +2615,8 @@ fn determine_decommission_final_state(items_failed: usize, was_cancelled: bool)
}
}
fn decommission_remaining_version_count(versions: &[rustfs_filemeta::FileInfo], expired: usize) -> usize {
versions
.iter()
.filter(|version| !version.tier_free_version())
.count()
.saturating_sub(expired)
fn decommission_remaining_version_count(total_versions: usize, expired: usize) -> usize {
total_versions.saturating_sub(expired)
}
fn should_skip_decommission_delete_marker(
@@ -2732,7 +2679,6 @@ fn decommission_remote_tiered_opts(
user_defined: version.metadata.clone(),
src_pool_idx,
data_movement: true,
incl_free_versions: version.tier_free_version(),
include_part_checksums: true,
http_preconditions: Some(crate::data_movement::data_movement_target_precondition()),
expected_bucket_incarnation_id,
@@ -3876,7 +3822,6 @@ impl ECStore {
let mut decommissioned: usize = 0;
let mut expired: usize = 0;
let mut free_version_disposition = DecommissionFreeVersionDisposition::default();
let mut cleanup_preflight_allowed_missing = Vec::new();
for version in fivs.versions.iter() {
@@ -3885,115 +3830,6 @@ impl ECStore {
}
decommission_cancel_signal_result(rx.is_cancelled())?;
if version.tier_free_version() {
let version_id = version.version_id.map(|v| v.to_string());
let mut migration_error = None;
let mut migrated = false;
let mut consumed = false;
let mut capacity_failure = false;
for _ in 0..3 {
match classify_decommission_free_version_attempt(
run_decommission_side_effect(&rx, &operation_gate, || async {
self.decommission_tiered_object(
bucket.as_str(),
&version.name,
version,
&decommission_remote_tiered_opts(
version,
version_id.clone(),
idx,
expected_bucket_incarnation_id,
),
)
.await
})
.await,
) {
DecommissionFreeVersionAttempt::Migrated => {
migrated = true;
migration_error = None;
break;
}
DecommissionFreeVersionAttempt::Consumed => {
consumed = true;
migration_error = None;
break;
}
DecommissionFreeVersionAttempt::CapacityFailure(err) => {
capacity_failure = true;
migration_error = Some(err);
break;
}
DecommissionFreeVersionAttempt::Retry(err) => migration_error = Some(err),
}
}
{
let mut pool_meta = self.pool_meta.write().await;
ensure_decommission_generation(&pool_meta, idx, generation)?;
if let Err(err) = count_decommission_item(&mut pool_meta, idx, 0, !migrated && !consumed) {
return Err(with_decommission_entry_context(
"count_decommission_item",
bucket.as_str(),
entry.name.as_str(),
err,
));
}
}
if migrated || consumed {
decommissioned += 1;
cleanup_preflight_allowed_missing.push(data_movement::source_cleanup_version_identity(version));
}
if migrated {
free_version_disposition.record_migrated();
} else if consumed {
free_version_disposition.record_consumed();
} else {
free_version_disposition.record_retained();
}
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %version.name,
version_id = ?version_id,
result = ?migration_error,
reason = if migrated {
DECOMMISSION_FREE_VERSION_MIGRATED_REASON
} else if consumed {
DECOMMISSION_FREE_VERSION_CONSUMED_REASON
} else {
DECOMMISSION_FREE_VERSION_RETAINED_REASON
},
state = if migrated {
"free_version_migrated"
} else if consumed {
"free_version_consumed"
} else {
"free_version_retained"
},
"Decommission free-version disposition recorded"
);
if capacity_failure {
return Err(with_decommission_entry_context(
"decommission_tier_free_version",
bucket.as_str(),
version.name.as_str(),
migration_error.expect("capacity failure must retain its error"),
));
}
if !migrated && !consumed {
break;
}
continue;
}
if run_decommission_side_effect(&rx, &operation_gate, || async {
should_skip_lifecycle_for_data_movement(
self.clone(),
@@ -4014,7 +3850,7 @@ impl ECStore {
continue;
}
let remaining_versions = decommission_remaining_version_count(&fivs.versions, expired);
let remaining_versions = decommission_remaining_version_count(fivs.versions.len(), expired);
if should_skip_decommission_delete_marker(version, remaining_versions, replication_config.is_some()) {
//
decommissioned += 1;
@@ -4293,24 +4129,6 @@ impl ECStore {
}
}
if free_version_disposition.total() > 0 {
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %entry.name,
free_versions_migrated = free_version_disposition.migrated,
free_versions_consumed = free_version_disposition.consumed,
free_versions_retained = free_version_disposition.retained,
free_versions_total = free_version_disposition.total(),
reason = DECOMMISSION_FREE_VERSION_DISPOSITION_REASON,
state = "free_version_disposition",
"Decommission free-version disposition summary"
);
}
if should_cleanup_decommission_source_entry(decommissioned, fivs.versions.len(), expired) {
if bucket_incarnation_fence.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
return Err(Error::other("decommission bucket incarnation fence was lost before source cleanup"));
@@ -4480,39 +4298,6 @@ impl ECStore {
.await
}
#[cfg(test)]
pub(crate) async fn decommission_entry_for_test_with_bucket_incarnation(
self: &Arc<Self>,
idx: usize,
entry: MetaCacheEntry,
bucket: String,
set: Arc<SetDisks>,
) -> Result<()> {
let expected_bucket_incarnation_id = if is_meta_bucketname(&bucket) {
None
} else {
Some(self.bucket_incarnation_id_from_disk(&bucket).await?)
};
self.decommission_entry(
CancellationToken::new(),
idx,
OffsetDateTime::now_utc(),
entry,
bucket,
set,
None,
None,
None,
expected_bucket_incarnation_id,
)
.await
}
#[cfg(test)]
pub(crate) async fn check_after_decommission_for_test(self: &Arc<Self>, idx: usize) -> Result<()> {
self.check_after_decommission(idx).await
}
#[tracing::instrument(skip(self, rx))]
async fn decommission_pool(
self: &Arc<Self>,
@@ -5377,7 +5162,6 @@ impl ECStore {
let lifecycle_config_cb = lifecycle_config.clone();
let object_lock_config_cb = object_lock_config.clone();
let store = Arc::clone(self);
let set_cb = Arc::clone(set);
let callback_rx_cb = callback_rx.clone();
let callback: ListCallback = Arc::new(move |entry: MetaCacheEntry| {
@@ -5387,7 +5171,6 @@ impl ECStore {
let lifecycle_config = lifecycle_config_cb.clone();
let object_lock_config = object_lock_config_cb.clone();
let store = Arc::clone(&store);
let set = Arc::clone(&set_cb);
let callback_rx = callback_rx_cb.clone();
Box::pin(async move {
if callback_rx.is_cancelled() {
@@ -5402,14 +5185,11 @@ impl ECStore {
return;
}
let fivs = match load_decommission_entry_exact_versions(
&set,
let fivs = match load_decommission_entry_versions(
&entry,
&bucket_name,
"check_after_decommission.file_info_versions",
)
.await
{
) {
Ok(fivs) => fivs,
Err(err) => {
let mut first_err = entry_error.lock().await;
@@ -5422,23 +5202,7 @@ impl ECStore {
};
let mut remaining = 0;
for version in fivs.versions.iter().chain(fivs.free_versions.iter()) {
if version.tier_free_version() {
remaining += 1;
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket_name,
object = %entry.name,
version_id = ?version.version_id,
reason = DECOMMISSION_FREE_VERSION_SWEEP_REASON,
state = "free_version_retained",
"Decommission final sweep retained a free version"
);
continue;
}
for version in &fivs.versions {
if version.deleted {
continue;
}
@@ -5553,6 +5317,13 @@ mod tests {
assert_eq!(determine_decommission_final_state(0, true), DecommissionFinalState::Failed);
}
#[test]
fn decommission_remaining_version_count_excludes_only_expired_versions() {
assert_eq!(decommission_remaining_version_count(1, 0), 1);
assert_eq!(decommission_remaining_version_count(2, 1), 1);
assert_eq!(decommission_remaining_version_count(1, 1), 0);
}
#[test]
fn lifecycle_action_removes_data_movement_version_rejects_delete_marker_action() {
assert!(!lifecycle_action_removes_data_movement_version(IlmAction::DeleteAction));
@@ -5609,21 +5380,6 @@ mod tests {
)));
}
#[test]
fn decommission_free_version_attempt_treats_missing_source_as_consumed() {
let attempt =
classify_decommission_free_version_attempt(Err(Error::ObjectNotFound("bucket".to_string(), "object".to_string())));
assert!(matches!(attempt, DecommissionFreeVersionAttempt::Consumed));
}
#[test]
fn decommission_free_version_attempt_preserves_capacity_failure() {
let attempt = classify_decommission_free_version_attempt(Err(Error::DiskFull));
assert!(matches!(attempt, DecommissionFreeVersionAttempt::CapacityFailure(Error::DiskFull)));
}
#[test]
fn decommission_delete_marker_copy_error_rejects_data_movement_overwrite() {
let err = Error::DataMovementOverwriteErr("bucket".to_string(), "object".to_string(), "version".to_string());
@@ -5805,12 +5561,6 @@ mod tests {
assert!(opts.include_part_checksums);
assert!(opts.http_preconditions.is_some());
assert_eq!(opts.expected_bucket_incarnation_id, Some(incarnation));
assert!(!opts.incl_free_versions);
let mut free_version = version;
free_version.set_tier_free_version();
let free_opts = decommission_remote_tiered_opts(&free_version, Some("free-version-id".to_string()), 9, Some(incarnation));
assert!(free_opts.incl_free_versions);
}
#[test]
-13
View File
@@ -1950,19 +1950,6 @@ mod tests {
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test]
fn test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
let mut free_version = cleanup_test_file_info("object.txt", Uuid::from_u128(2), "tier-cleanup");
free_version.deleted = true;
free_version.set_tier_free_version();
let expected = cleanup_test_versions(vec![migrated.clone(), free_version.clone()]);
let current = cleanup_test_versions(vec![migrated]);
let allowed_missing = vec![source_cleanup_version_identity(&free_version)];
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test]
fn test_decommission_cleanup_preflight_rejects_unexpected_missing_version() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
+76 -6
View File
@@ -8000,10 +8000,15 @@ impl DiskAPI for LocalDisk {
use std::io::Write as _;
let file_path = self.io_get_object_path(volume, path)?;
let lock_path = file_path.with_extension("rustfs-cas.lock");
let path = path.to_string();
let sync_metadata = effective_durability(volume).syncs_commit_metadata();
return Ok(tokio::task::spawn_blocking(move || {
// A persistent directory lock bounds metadata growth. Removing
// per-target lock files can split flock ownership across inodes.
let lock_path = file_path
.parent()
.ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))?
.join(".rustfs-cas.lock");
let lock = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
@@ -8068,7 +8073,25 @@ impl DiskAPI for LocalDisk {
.map_err(DiskError::from)??);
}
#[cfg(not(unix))]
#[cfg(windows)]
{
let file_path = self.io_get_object_path(volume, path)?;
let sync_metadata = effective_durability(volume).syncs_commit_metadata();
let publication_root = self.publication_root.clone();
return Ok(tokio::task::spawn_blocking(move || {
os::compare_and_update_control_file(
&file_path,
expected.as_deref(),
replacement.as_deref(),
sync_metadata,
&publication_root,
)
})
.await
.map_err(DiskError::from)??);
}
#[cfg(not(any(unix, windows)))]
{
let _ = (volume, path, expected, replacement);
Err(DiskError::MethodNotAllowed)
@@ -21823,9 +21846,9 @@ mod test {
assert!(matches!(results[1].as_ref().unwrap_err(), DiskError::Io(_)));
}
#[cfg(unix)]
#[cfg(any(unix, windows))]
#[tokio::test]
async fn conditional_file_update_never_deletes_a_new_owner() {
async fn windows_and_unix_conditional_file_update_never_deletes_a_new_owner() {
use tempfile::tempdir;
let dir = tempdir().expect("temp dir should be created");
@@ -21856,8 +21879,18 @@ mod test {
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.await
.expect("new owner marker should remain"),
owner_b
owner_b.clone()
);
assert_eq!(
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, Some(owner_b), None)
.await
.expect("current owner should remove marker"),
ConditionalFileUpdate::Updated
);
assert!(matches!(
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH).await,
Err(DiskError::FileNotFound)
));
}
#[cfg(unix)]
@@ -21872,7 +21905,10 @@ mod test {
let marker_path = disk
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.expect("marker path should resolve");
let lock_path = marker_path.with_extension("rustfs-cas.lock");
let lock_path = marker_path
.parent()
.expect("marker path should have a parent")
.join(".rustfs-cas.lock");
let lock = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
@@ -21893,6 +21929,40 @@ mod test {
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock));
}
#[cfg(windows)]
#[tokio::test]
async fn windows_conditional_file_update_returns_would_block_when_marker_lock_is_contended() {
let dir = tempfile::tempdir().expect("temp dir should be created");
let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse");
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
ensure_test_volume(&disk, RUSTFS_META_BUCKET).await;
let marker_path = disk
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.expect("marker path should resolve");
let lock_path = marker_path
.parent()
.expect("marker path should have a parent")
.join(".rustfs-cas.lock");
let lock = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(lock_path)
.expect("marker lock should open");
lock.try_lock().expect("marker lock should be held");
let err = tokio::time::timeout(
Duration::from_secs(1),
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, None, Some(Bytes::from_static(b"owner"))),
)
.await
.expect("contended conditional update must not block")
.expect_err("contended conditional update must retry");
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock));
}
#[cfg(target_os = "linux")]
#[tokio::test]
async fn replacement_io_paths_stay_under_the_mount_lease() {
+85
View File
@@ -12,6 +12,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#[cfg(windows)]
use crate::disk::ConditionalFileUpdate;
use crate::disk::error::DiskError;
use crate::disk::error::Result;
use crate::disk::error_conv::to_file_error;
@@ -3458,6 +3460,89 @@ fn read_windows_relative_file(file_path: &Path, parent_guard: &ExistingBaseDirec
Ok(Some(data))
}
#[cfg(windows)]
pub(crate) fn compare_and_update_control_file(
file_path: &Path,
expected: Option<&[u8]>,
replacement: Option<&[u8]>,
sync_metadata: bool,
publication_root: &PublicationRoot,
) -> io::Result<ConditionalFileUpdate> {
use windows_sys::{
Wdk::Storage::FileSystem::{
FILE_NON_DIRECTORY_FILE, FILE_OPEN, FILE_OPEN_IF, FILE_OPEN_REPARSE_POINT, FILE_SYNCHRONOUS_IO_NONALERT,
},
Win32::Storage::FileSystem::{
DELETE, FILE_ATTRIBUTE_NORMAL, FILE_READ_ATTRIBUTES, FILE_SHARE_READ, FILE_SHARE_WRITE, FILE_WRITE_DATA, SYNCHRONIZE,
},
};
let parent = file_path
.parent()
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file has no parent"))?;
let parent_guard = lock_windows_directory_tree(parent, Some(parent), publication_root)?;
let lock = open_windows_relative(
parent_guard.last_handle()?,
std::ffi::OsStr::new(".rustfs-cas.lock"),
SYNCHRONIZE | FILE_READ_ATTRIBUTES | FILE_WRITE_DATA,
FILE_SHARE_READ | FILE_SHARE_WRITE,
FILE_OPEN_IF,
FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT,
FILE_ATTRIBUTE_NORMAL,
true,
)?;
validate_windows_owned_file(&lock)?;
match lock.as_file().try_lock() {
Ok(()) => {}
Err(std::fs::TryLockError::WouldBlock) => return Err(io::Error::from(io::ErrorKind::WouldBlock)),
Err(std::fs::TryLockError::Error(err)) => return Err(err),
}
let current = read_windows_relative_file(file_path, &parent_guard)?;
let matches = match (&current, expected) {
(None, None) => true,
(Some(current), Some(expected)) => current.as_slice() == expected,
_ => false,
};
if !matches {
return Ok(match current {
None => ConditionalFileUpdate::Missing,
Some(_) => ConditionalFileUpdate::Mismatch,
});
}
match replacement {
Some(replacement) => RenameDestinationPathGuard {
directory: parent.to_path_buf(),
_directory_guard: parent_guard,
}
.write_file_for_path_access(file_path, replacement, sync_metadata, sync_metadata)?,
None => {
let file_name = file_path
.file_name()
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file must have a name"))?;
let file = open_windows_relative(
parent_guard.last_handle()?,
file_name,
DELETE | SYNCHRONIZE | FILE_READ_ATTRIBUTES,
FILE_SHARE_READ,
FILE_OPEN,
FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT,
0,
true,
)?;
validate_windows_owned_file(&file)?;
set_windows_file_delete_on_close(&file, true)?;
drop(file);
if sync_metadata {
fsync_dir_std(parent)?;
}
}
}
Ok(ConditionalFileUpdate::Updated)
}
#[cfg(windows)]
fn open_windows_directory_component(
parent: &WindowsDirectoryHandle,
+6 -38
View File
@@ -784,24 +784,6 @@ pub(crate) fn create_deferred_bitrot_reader_with_stripe_handle(
///
/// # Returns
/// A Result containing the BitrotWriterWrapper or an error
/// Size hint handed to `DiskAPI::create_file` for a bitrot-wrapped shard.
///
/// A known length is grown by one checksum per shard so the on-disk file size
/// matches what the bitrot writer emits. A negative length is the
/// unknown-size sentinel (`HashReader::SIZE_PRESERVE_LAYER`, used by SSE and
/// compression) and must be preserved: `RemoteDisk::create_file` forwards it
/// in the `put_file_stream` query, and the receiver only treats `size > 0` as
/// a fixed body length when locating the authenticated trailer. Clamping it
/// to `0` would claim an empty body and misframe the stream. `0` stays `0`
/// because a genuinely empty object still means an empty body.
fn bitrot_create_file_size(length: i64, shard_size: usize, checksum_algo: &HashAlgorithm) -> i64 {
if length <= 0 {
return length;
}
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
}
pub async fn create_bitrot_writer(
is_inline_buffer: bool,
disk: Option<&DiskStore>,
@@ -814,7 +796,12 @@ pub async fn create_bitrot_writer(
let writer = if is_inline_buffer {
CustomWriter::new_inline_buffer()
} else if let Some(disk) = disk {
let length = bitrot_create_file_size(length, shard_size, &checksum_algo);
let length = if length > 0 {
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
} else {
0
};
let file = disk.create_file("", volume, path, length).await?;
#[cfg(feature = "hotpath")]
@@ -833,25 +820,6 @@ mod tests {
use rustfs_rio::ChunkReader;
use std::collections::VecDeque;
#[test]
fn bitrot_create_file_size_grows_known_length_by_checksums() {
// 10 bytes over 4-byte shards = 3 shards, each followed by a 32-byte hash.
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::HighwayHash256), 10 + 3 * 32);
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::None), 10);
}
#[test]
fn bitrot_create_file_size_keeps_empty_and_unknown_distinct() {
assert_eq!(bitrot_create_file_size(0, 4, &HashAlgorithm::HighwayHash256), 0);
// SSE/compression streams advertise SIZE_PRESERVE_LAYER (-1); the remote
// put_file_stream receiver relies on a non-positive size to parse the auth
// trailer from the stream tail, so the sentinel must survive untouched.
assert_eq!(
bitrot_create_file_size(rustfs_rio::HashReader::SIZE_PRESERVE_LAYER, 4, &HashAlgorithm::HighwayHash256),
rustfs_rio::HashReader::SIZE_PRESERVE_LAYER
);
}
struct TestChunkReader {
chunks: VecDeque<Bytes>,
}
-297
View File
@@ -227,70 +227,6 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap<Strin
restore_operation_id_from_metadata(metadata)
}
async fn inspect_decommission_tier_free_version_target(
disk: &DiskStore,
bucket: &str,
object: &str,
source: &FileInfo,
) -> Result<bool> {
let raw = match disk.read_xl(bucket, object, false).await {
Ok(raw) => raw,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(false),
Err(err) => return Err(err.into()),
};
let meta = FileMeta::load(&raw.buf)?;
let source_version_id = source.version_id.filter(|version_id| !version_id.is_nil());
let mut matching_count = 0;
let mut all_matching_versions_equivalent = true;
for existing in meta
.versions
.iter()
.filter(|version| version.header.version_id.filter(|version_id| !version_id.is_nil()) == source_version_id)
{
matching_count += 1;
let existing = existing.into_fileinfo(bucket, object, true)?;
existing.validate_for_metadata_read()?;
if !existing.tier_free_version() || !crate::store::tiered_data_movement_source_matches(source, &existing)? {
all_matching_versions_equivalent = false;
}
}
if matching_count == 0 {
return Ok(false);
}
if matching_count == 1 && all_matching_versions_equivalent {
return Ok(true);
}
Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
source_version_id.map(|version_id| version_id.to_string()).unwrap_or_default(),
)
.into())
}
fn ensure_decommission_tier_free_version_commit_fence(bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> {
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
|| opts
.bucket_lifecycle_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
Ok(())
}
impl SetDisks {
pub(super) async fn require_current_restore_operation_id(
&self,
@@ -4729,98 +4665,6 @@ fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_
}
impl SetDisks {
/// Publish an internal tier free-version record without changing its
/// delete-marker shape or remote-tier identity. The caller holds the
/// source and target object locks; a write quorum is required before the
/// source cleanup may remove the original record.
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tier_free_version(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
if !fi.deleted || !fi.tier_free_version() {
return Err(Error::other("decommission tier free-version write requires a free version record"));
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
self.validate_decommission_tier_free_version_target(bucket, object, fi)
.await?;
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let write_quorum = self.default_write_quorum();
let futures = disks.into_iter().map(|disk| {
let file_info = fi.clone();
async move {
if let Some(disk) = disk {
disk.write_metadata("", bucket, object, file_info).await
} else {
Err(DiskError::DiskNotFound)
}
}
});
let mut errs = Vec::new();
for result in join_all(futures).await {
match result {
Ok(_) => errs.push(None),
Err(err) => errs.push(Some(err)),
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object)
}
pub(crate) async fn validate_decommission_tier_free_version_target(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
) -> Result<()> {
// The caller holds the source and target object locks. Inspect every
// target disk before an idempotent return or metadata fan-out so a
// sub-quorum conflict cannot be hidden by a successful quorum.
let disks = self.disks.read().await.clone();
let preflight = disks
.iter()
.flatten()
.map(|disk| inspect_decommission_tier_free_version_target(disk, bucket, object, fi));
for result in join_all(preflight).await {
result?;
}
Ok(())
}
pub(crate) async fn has_decommission_tier_free_version_write_quorum(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<bool> {
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let preflight = disks.iter().map(|disk| async {
match disk {
Some(disk) => inspect_decommission_tier_free_version_target(disk, bucket, object, fi).await,
None => Ok(false),
}
});
let mut equivalent = 0;
for result in join_all(preflight).await {
if result? {
equivalent += 1;
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
Ok(equivalent >= self.default_write_quorum())
}
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tiered_object(
&self,
@@ -10043,147 +9887,6 @@ mod tests {
assert_ne!(updated.erasure.distribution, original.erasure.distribution);
}
#[tokio::test]
async fn decommission_tier_free_version_preserves_remote_identity() {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "free-version-decommission";
let object = "object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("target bucket should exist before free-version migration");
let version_id = Uuid::new_v4();
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(version_id),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
// Decoded free versions always carry the on-disk free-version
// suffix alongside the in-memory tier marker; mirror that here so
// the record satisfies delete-marker metadata validation.
rustfs_utils::http::metadata_compat::insert_str(
&mut free_version.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_FREE_VERSION,
String::new(),
);
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("free-version metadata should reach the target quorum");
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("replaying the same free-version metadata should be idempotent");
let versions = set_disks
.load_file_info_versions_exact(bucket, object)
.await
.expect("migrated free-version metadata should decode")
.expect("migrated free-version metadata should exist");
let migrated = versions
.versions
.iter()
.find(|version| version.version_id == Some(version_id))
.expect("free version should be present on the target");
assert_eq!(
versions
.versions
.iter()
.filter(|version| version.version_id == Some(version_id))
.count(),
1
);
assert!(migrated.tier_free_version());
assert_eq!(migrated.transition_tier, "WARM-TIER");
assert_eq!(migrated.transitioned_objname, "remote/object");
}
#[tokio::test]
async fn decommission_tier_free_version_resume_requires_write_quorum() {
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
let bucket = "free-version-decommission-resume";
let object = "object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("target bucket should exist before free-version migration");
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(Uuid::new_v4()),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
// Decoded free versions always carry the on-disk free-version
// suffix alongside the in-memory tier marker; mirror that here so
// the record satisfies delete-marker metadata validation.
rustfs_utils::http::metadata_compat::insert_str(
&mut free_version.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_FREE_VERSION,
String::new(),
);
let opts = ObjectOptions::default();
let disks = set_disks.get_disks_internal().await;
for disk in disks.iter().take(2).flatten() {
disk.write_metadata("", bucket, object, free_version.clone())
.await
.expect("partial first attempt should leave equivalent metadata");
}
assert!(
!set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("partial target metadata should remain valid"),
"write-quorum-minus-one must not be accepted as an idempotent migration"
);
disks[2]
.as_ref()
.expect("third target disk should be online")
.write_metadata("", bucket, object, free_version.clone())
.await
.expect("third equivalent target write should complete quorum");
assert!(
set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("write-quorum target metadata should remain valid")
);
}
#[test]
fn decommission_tier_free_version_commit_rejects_lost_fence() {
let opts = ObjectOptions {
namespace_lock_fence: Some(NamespaceLockFence::lost_for_test()),
..Default::default()
};
let err = ensure_decommission_tier_free_version_commit_fence("bucket", "object", &opts)
.expect_err("lost target lock must fail the free-version commit");
assert!(matches!(
err,
Error::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
required: 1,
achieved: 0,
..
}
));
}
#[test]
fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() {
let errs = vec![None, None, Some(DiskError::DiskNotFound), None];
-450
View File
@@ -553,8 +553,6 @@ mod tests {
should_retry_format_load, should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay,
};
#[cfg(feature = "test-util")]
use crate::disk::DiskAPI;
#[cfg(feature = "test-util")]
use crate::{
bucket::lifecycle::{
lifecycle::{TRANSITION_PENDING, TransitionOptions},
@@ -1309,77 +1307,6 @@ mod tests {
});
}
#[cfg(feature = "test-util")]
async fn seed_transitioned_free_version(
ctx: &Arc<crate::runtime::instance::InstanceContext>,
store: &Arc<crate::store::ECStore>,
bucket: &str,
object: &str,
) -> (uuid::Uuid, uuid::Uuid) {
let tier_name = format!("DECOMFREE{}", uuid::Uuid::new_v4().simple());
register_mock_tier(&ctx.tier_config_mgr(), &tier_name).await;
let mut reader = PutObjReader::from_vec(b"transitioned source bytes".to_vec());
let source = store.pools[0]
.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write transitioned decommission source");
let source_version = source.version_id.expect("transitioned source must be versioned");
store.pools[0]
.transition_object(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: Some(source_version.to_string()),
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name,
etag: source.etag.clone().expect("transitioned source must have an ETag"),
..Default::default()
},
mod_time: source.mod_time,
..Default::default()
},
)
.await
.expect("transition source before decommission");
store.pools[0]
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(source_version.to_string()),
..Default::default()
},
)
.await
.expect("delete transitioned source version");
let versions = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("source versions should decode after transition delete")
.expect("source free version should remain after transition delete");
let free_version = versions
.versions
.iter()
.find(|version| version.tier_free_version())
.and_then(|version| version.version_id)
.expect("transition delete should create a free version");
(source_version, free_version)
}
async fn write_decommission_test_multipart_source(
store: &Arc<crate::store::ECStore>,
pool_idx: usize,
@@ -3919,383 +3846,6 @@ mod tests {
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_skips_cleanup_only_marker_when_free_version_is_present() {
let temp_dir = tempfile::tempdir().expect("create free-version decommission store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-marker", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-marker-{}", uuid::Uuid::new_v4());
let object = "free-marker-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create free-version decommission bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free-version metadata should decode")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
})
.expect("source free-version identity should be present before decommission");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target_w = store.pools[1]
.put_object(
&bucket,
object,
&mut target_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write unrelated target version");
let target_w_version = target_w.version_id.expect("target version should have an id");
let marker = store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write cleanup-only delete marker");
assert!(marker.delete_marker);
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.clone(),
source_set.clone(),
)
.await
.expect("real decommission entry should migrate the free version");
let target_versions = store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("target versions should decode")
.expect("target free version should be present");
assert!(
target_versions
.versions
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let migrated_free = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
.expect("migrated free-version identity should remain readable from target disks");
assert!(
crate::store::tiered_data_movement_source_matches(&source_free, migrated_free)
.expect("migrated free-version identity should decode")
);
let retained_w = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(target_w_version))
.expect("unrelated target version should remain");
assert!(!retained_w.deleted && !retained_w.tier_free_version());
assert_eq!(retained_w.size, target_w.size);
assert_eq!(retained_w.get_etag(), target_w.etag);
assert!(
target_versions
.versions
.iter()
.all(|version| { version.tier_free_version() || !version.deleted })
);
assert!(
source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source versions should be readable after cleanup")
.is_none(),
"successful free-version migration should permit source cleanup"
);
let (heal_versions, _, _) = store
.heal_walk_versions_page(1, 0, &bucket, "", None, 2, 16, true)
.await
.expect("heal walk should decode the migrated free version");
let free_version_string = free_version.to_string();
let healed_free = heal_versions
.iter()
.find(|version| version.version_id.as_deref() == Some(free_version_string.as_str()))
.expect("heal walk should surface the migrated free version");
let healed_info = healed_free
.lifecycle_object_info
.as_ref()
.expect("heal walk should retain lifecycle identity for the migrated free version");
assert!(healed_info.transitioned_object.free_version);
assert_eq!(healed_info.transitioned_object.tier, source_free.transition_tier);
assert_eq!(healed_info.transitioned_object.name, source_free.transitioned_objname);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_allows_free_version_consumed_before_source_lock() {
let temp_dir = tempfile::tempdir().expect("create consumed free-version store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-consumed", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-consumed-{}", uuid::Uuid::new_v4());
let object = "free-consumed-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create consumed free-version bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
let barrier = crate::store::object::DecommissionFreeVersionSourceRaceBarrier::install(&bucket, object);
let decommission = tokio::spawn({
let store = store.clone();
let bucket = bucket.clone();
let source_set = source_set.clone();
async move {
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket,
source_set,
)
.await
}
});
barrier.wait_until_paused().await;
store.pools[0]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(free_version.to_string()),
incl_free_versions: true,
..Default::default()
},
)
.await
.expect("lifecycle should consume the source free version before decommission locks it");
assert!(
source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("consumed source metadata should remain readable")
.is_none(),
"the lifecycle delete should remove the source free version"
);
barrier.release();
decommission
.await
.expect("decommission task should join")
.expect("a concurrently consumed free version should not fail source cleanup");
assert!(
store.pools[1]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("target metadata should remain readable")
.is_none(),
"an already consumed free version should not be recreated on the target"
);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source() {
let temp_dir = tempfile::tempdir().expect("create sub-quorum free-version store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-conflict", &[4, 4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("decom-free-conflict-{}", uuid::Uuid::new_v4());
let object = "free-conflict-object";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create sub-quorum conflict bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free version should decode before crash replay setup")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version))
})
.expect("source free version should be available for crash replay setup");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target = store.pools[1]
.put_object(
&bucket,
object,
&mut target_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed ordinary target version");
let target_version = target.version_id.expect("target version must have an ID");
let target_disks = store.pools[1].get_disks_by_key(object).disks.read().await.clone();
for disk in target_disks.iter().skip(1) {
disk.as_ref()
.expect("target crash replay quorum disk should be online")
.write_metadata("", &bucket, object, source_free.clone())
.await
.expect("seed an equivalent free version on the target quorum");
}
let conflict_path = temp_dir
.path()
.join(format!("pool1/set0/disk0/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let encoded = tokio::fs::read(&conflict_path)
.await
.expect("target metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("target metadata should decode");
let target_index = metadata
.versions
.iter()
.position(|version| version.header.version_id == Some(target_version))
.expect("target version should be present on the conflict disk");
let mut target_meta = metadata.versions[target_index]
.parse_version_meta()
.expect("target version metadata should decode");
target_meta
.object
.as_mut()
.expect("target conflict must remain an ordinary object")
.version_id = Some(free_version);
metadata.versions[target_index] = target_meta.try_into().expect("conflict metadata should encode");
let expected_conflict_meta = metadata.versions[target_index].meta.clone();
let expected_conflict = metadata.versions[target_index]
.into_fileinfo(&bucket, object, true)
.expect("conflict metadata should decode as an ordinary object");
let duplicate_free: rustfs_filemeta::FileMetaShallowVersion = rustfs_filemeta::FileMetaVersion::from(source_free.clone())
.try_into()
.expect("duplicate free metadata should encode");
metadata.versions.insert(target_index, duplicate_free);
tokio::fs::write(&conflict_path, metadata.marshal_msg().expect("conflict metadata should encode"))
.await
.expect("write sub-quorum conflict metadata");
mark_test_pool_decommissioning(&store, 0).await;
let source_set = store.pools[0].get_disks_by_key(object);
store
.decommission_entry_for_test_with_bucket_incarnation(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.clone(),
source_set.clone(),
)
.await
.expect("conflicted decommission entry should retain the source and retry later");
let source_versions = source_set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("retained source versions should decode")
.expect("source free version should be retained after conflict");
assert!(
source_versions
.versions
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let post_encoded = tokio::fs::read(&conflict_path)
.await
.expect("conflict metadata should remain readable");
let post_metadata = FileMeta::load(&post_encoded).expect("post-conflict metadata should decode");
let same_id = post_metadata
.versions
.iter()
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
assert_eq!(same_id.len(), 2, "conflict metadata should retain both same-ID records");
assert_eq!(same_id.iter().filter(|version| version.header.free_version()).count(), 1);
assert_eq!(same_id.iter().filter(|version| !version.header.free_version()).count(), 1);
let post_conflict = same_id
.into_iter()
.find(|version| !version.header.free_version())
.expect("ordinary conflict version must remain addressable by the source ID");
let post_conflict_info = post_conflict
.into_fileinfo(&bucket, object, true)
.expect("post-conflict ordinary metadata should decode");
assert!(!post_conflict_info.deleted && !post_conflict_info.tier_free_version());
assert_eq!(post_conflict.meta, expected_conflict_meta);
assert_eq!(post_conflict_info.size, expected_conflict.size);
assert_eq!(post_conflict_info.data_dir, expected_conflict.data_dir);
assert_eq!(post_conflict_info.metadata, expected_conflict.metadata);
assert_eq!(post_conflict_info.get_etag(), expected_conflict.get_etag());
for disk_index in 0..4 {
let target_path = temp_dir
.path()
.join(format!("pool1/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let target_encoded = tokio::fs::read(&target_path)
.await
.expect("target metadata should remain readable");
let target_meta = FileMeta::load(&target_encoded).expect("target metadata should decode");
let same_id = target_meta
.versions
.iter()
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
if disk_index == 0 {
assert_eq!(same_id.len(), 2);
assert!(same_id[0].header.free_version());
assert!(!same_id[1].header.free_version());
} else {
assert_eq!(same_id.len(), 1);
assert!(same_id[0].header.free_version());
}
}
let sweep_err = store
.check_after_decommission_for_test(0)
.await
.expect_err("final sweep must report the retained free version");
assert!(
sweep_err.to_string().contains("version(s) were found"),
"unexpected final sweep error: {sweep_err}"
);
shutdown.cancel();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn versioned_batch_delete_marker_skips_decommission_source() {
+1 -1
View File
@@ -151,7 +151,7 @@ pub(crate) mod init_format;
pub(crate) mod list_objects;
mod multipart;
mod object;
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence, tiered_data_movement_source_matches};
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence};
pub use object::{
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
SnapshotConsistencyError,
+15 -143
View File
@@ -490,82 +490,6 @@ fn decommission_mutation_fence_for_test(
.map(|hook| hook.fence.clone())
}
#[cfg(test)]
struct DecommissionFreeVersionSourceRaceState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct DecommissionFreeVersionSourceRaceBarrier {
state: Arc<DecommissionFreeVersionSourceRaceState>,
}
#[cfg(test)]
static DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<DecommissionFreeVersionSourceRaceState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl DecommissionFreeVersionSourceRaceBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(DecommissionFreeVersionSourceRaceState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison");
assert!(slot.is_none(), "decommission free-version source race 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("decommission should pause before acquiring the free-version source lock");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for DecommissionFreeVersionSourceRaceBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_decommission_free_version_before_source_lock(bucket: &str, object: &str) {
let state = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission free-version source race barrier should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.object == object)
.cloned();
if let Some(state) = state {
state.arrived.notify_one();
state.release.notified().await;
}
}
pub(crate) struct SourceCleanupMutationFence {
guard: ObjectLockDiagGuard,
source_lock_covered: bool,
@@ -1570,15 +1494,13 @@ fn is_equivalent_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo,
&& source_actual_size == target_actual_size
}
pub(crate) fn tiered_data_movement_source_matches(
fn tiered_data_movement_source_matches(
expected: &rustfs_filemeta::FileInfo,
current: &rustfs_filemeta::FileInfo,
) -> Result<bool> {
let expected_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&expected.metadata)?;
let current_backend = crate::services::tier::tier::tier_destination_id_from_metadata(&current.metadata)?;
Ok(expected.version_id == current.version_id
&& expected.deleted == current.deleted
&& expected.tier_free_version() == current.tier_free_version()
&& expected.data_dir == current.data_dir
&& expected.mod_time == current.mod_time
&& expected.size == current.size
@@ -1592,15 +1514,6 @@ pub(crate) fn tiered_data_movement_source_matches(
&& expected_backend == current_backend)
}
fn decommission_free_version_overwrite_error(bucket: &str, object: &str, version_id: Option<Uuid>) -> Error {
StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
version_id.map(|id| id.to_string()).unwrap_or_default(),
)
.into()
}
fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool {
target_pool_idx != src_pool_idx
}
@@ -2308,23 +2221,6 @@ impl ECStore {
)
}
async fn has_equivalent_data_movement_tier_free_version(
&self,
bucket: &str,
object: &str,
source: &rustfs_filemeta::FileInfo,
opts: &ObjectOptions,
target_pool_idx: usize,
) -> Result<bool> {
let pool = self
.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement target pool {target_pool_idx}")))?;
pool.get_disks_by_key(object)
.has_decommission_tier_free_version_write_quorum(bucket, object, source, opts)
.await
}
fn resolve_decommission_target_pool_idx_result(result: Result<usize>, bucket: &str, object: &str) -> Result<usize> {
result.map_err(|err| Error::other(format!("failed to select decommission target pool for {bucket}/{object}: {err}")))
}
@@ -2344,10 +2240,6 @@ impl ECStore {
check_put_object_args(bucket, object)?;
let mut opts = opts.clone();
let is_free_version = fi.tier_free_version();
if is_free_version {
opts.incl_free_versions = true;
}
let bucket_incarnation_fence = if is_meta_bucketname(bucket) {
None
} else {
@@ -2385,10 +2277,6 @@ impl ECStore {
&object,
)?
};
#[cfg(test)]
if is_free_version {
pause_decommission_free_version_before_source_lock(bucket, logical_object).await;
}
let _object_guards = self
.acquire_data_movement_object_write_locks(bucket, &object, opts.src_pool_idx, idx, &mut opts)
.await?;
@@ -2406,7 +2294,7 @@ impl ECStore {
versions
.versions
.iter()
.find(|current| current.version_id == fi.version_id && current.tier_free_version() == is_free_version)
.find(|current| current.version_id == fi.version_id && !current.tier_free_version())
})
.ok_or_else(|| to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]))?;
if !tiered_data_movement_source_matches(fi, current_source)? {
@@ -2421,40 +2309,24 @@ impl ECStore {
.get_available_pool_idx_excluding(bucket, &object, fi.size, opts.src_pool_idx)
.await;
let target_pool_idx = resolve_data_movement_resume_target_pool(idx, resume_target_pool_idx, opts.src_pool_idx);
if is_free_version && target_pool_idx == opts.src_pool_idx {
return Err(Error::DiskFull);
}
let equivalent = if is_free_version {
self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, target_pool_idx)
.await?
} else {
self.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await?
};
if equivalent {
return Ok(());
}
return Err(decommission_free_version_overwrite_error(bucket, &object, fi.version_id));
}
let result = if is_free_version {
if self
.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, idx)
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await?
{
return Ok(());
}
self.pools[idx]
.get_disks_by_key(&object)
.decommission_tier_free_version(bucket, &object, &fi, &opts)
.await
} else {
self.pools[idx]
.get_disks_by_key(&object)
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await
};
return Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
opts.version_id.clone().unwrap_or_default(),
));
}
let result = self.pools[idx]
.get_disks_by_key(&object)
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await;
if matches!(result, Err(Error::PreconditionFailed)) {
if self
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, idx)
-19
View File
@@ -90,21 +90,6 @@ fn legacy_data_key_for_version(version_id: Option<Uuid>) -> Option<String> {
pub const TRANSITION_COMPLETE: &str = "complete";
pub const TRANSITION_PENDING: &str = "pending";
/// xl.meta key marking a tier free-version record.
///
/// A free version is a delete-marker-shaped cleanup hint appended by
/// [`MetaObject::delete_version`] when a version whose remote transition
/// completed is removed from xl.meta; it carries the remote tier identity for
/// an idempotent remote delete and is never a user-visible version
/// (`num_versions` excludes it). While the record exists it is consumed by the
/// lifecycle free-version recovery scan and the usage scanner, which re-enqueue
/// the pending remote delete, and by heal metadata walks. On S3 and lifecycle
/// delete paths the same obligation is also carried by a committed tier-journal
/// entry; deletes without such an entry (for example a removed version whose
/// transition state decodes as unknown) rely on this record alone until the
/// worker removes it after a successful remote delete. Decommission preserves
/// the record and its remote identity on the target pool before source cleanup
/// — see docs/architecture/decommission-compatibility.md.
pub const FREE_VERSION: &str = "free-version";
pub const TRANSITION_STATUS: &str = "transition-status";
@@ -462,10 +447,6 @@ impl FileMeta {
};
if let Some(fidx) = existing_idx {
let existing = self.versions[fidx].parse_version_meta()?;
if existing.free_version() != version.free_version() {
return Err(Error::other("cannot replace a free version with a non-free version"));
}
return self.set_idx(fidx, version);
}
-9
View File
@@ -2725,15 +2725,6 @@ impl MetaObject {
self.meta_sys.retain(|k, _| !k.starts_with("X-Amz-Restore"));
}
/// Builds the free-version cleanup record appended when a transitioned
/// version is removed from xl.meta. The record keeps the remote tier
/// identity so the lifecycle worker can issue the idempotent remote delete
/// and only then remove the record; until then the recovery scan and the
/// usage scanner keep re-enqueueing it. S3 and lifecycle deletes also
/// persist a committed tier-journal entry for the same remote delete. The
/// decommission path copies this record unchanged before source cleanup,
/// including when the transition state is unknown — see
/// docs/architecture/decommission-compatibility.md.
pub fn init_free_version(&self, fi: &FileInfo) -> Result<(FileMetaVersion, bool)> {
if fi.skip_tier_free_version() {
return Ok((FileMetaVersion::default(), false));
+1
View File
@@ -91,6 +91,7 @@ metrics = { workspace = true }
base64 = { workspace = true }
bytes = { workspace = true }
crc-fast = { workspace = true }
sha2 = { workspace = true }
[dev-dependencies]
serde_json = { workspace = true, features = ["raw_value"] }
+5
View File
@@ -373,6 +373,11 @@ impl ErasureSetHealer {
set_disk_id: &str,
buckets: &[String],
) -> Result<(ResumeManager, CheckpointManager)> {
if self.replacement_task_id.is_none() && CheckpointManager::is_blocked(&self.disk, task_id).await {
return Err(Error::TaskExecutionFailed {
message: format!("Resume task {task_id} has a blocked checkpoint"),
});
}
// check if resume state exists
let has_resume_state = if self.replacement_task_id.is_some() {
ResumeManager::has_replacement_intent(&self.disk, task_id).await
+1
View File
@@ -51,6 +51,7 @@ const RESUME_STATE_FILE: &str = "ahm_resume_state.json";
const REPLACEMENT_INTENT_FILE: &str = "ahm_replacement_intent.json";
const RESUME_PROGRESS_FILE: &str = "ahm_progress.json";
pub(super) const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json";
pub(super) const RESUME_CHECKPOINT_BLOCKED_FILE: &str = "ahm_checkpoint.blocked";
const REPLACEMENT_COMPLETION_PROOF_FILE: &str = "ahm_replacement_completion_proof.json";
const REPLACEMENT_RECOVERY_DIR: &str = "ahm-replacement";
const REPLACEMENT_INTENT_SEAL_FILE: &str = "ahm_replacement_intent_seal";
+317 -21
View File
@@ -13,26 +13,31 @@
// limitations under the License.
use crate::{Error, Result};
use base64::Engine as _;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::HashSet;
use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH};
use tokio::sync::RwLock;
use tokio::sync::{Mutex as AsyncMutex, RwLock};
use tracing::{debug, warn};
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
use super::super::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes};
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt, RUSTFS_META_BUCKET};
use super::{
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_FILE, delete_resume_file, path_to_str,
validate_resume_task_id,
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_BLOCKED_FILE, RESUME_CHECKPOINT_FILE,
delete_resume_file, path_to_str, validate_resume_task_id,
};
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
const RESUME_CHECKPOINT_DIGEST_FILE: &str = "ahm_checkpoint.sha256";
const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 5;
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
/// to the new `compose_key` identities, so a stale checkpoint is discarded.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6;
/// resume checkpoint
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -57,6 +62,11 @@ pub struct ResumeCheckpoint {
pub failed_objects: HashSet<String>,
/// skipped objects
pub skipped_objects: HashSet<String>,
/// Integrity digest over the checkpoint with this field set to `None`.
/// Keeping it in the checkpoint makes the payload and its authentication
/// record one CAS generation instead of two independently-written files.
#[serde(default)]
pub integrity_digest: Option<String>,
}
impl ResumeCheckpoint {
@@ -70,6 +80,7 @@ impl ResumeCheckpoint {
processed_objects: HashSet::new(),
failed_objects: HashSet::new(),
skipped_objects: HashSet::new(),
integrity_digest: None,
}
}
@@ -116,17 +127,111 @@ pub struct CheckpointManager {
disk: DiskStore,
checkpoint: Arc<RwLock<ResumeCheckpoint>>,
throttle: Mutex<PersistThrottle>,
save_lock: AsyncMutex<()>,
last_saved: Mutex<Option<EcstoreDiskBytes>>,
}
impl CheckpointManager {
fn blocked_path(task_id: &str) -> std::path::PathBuf {
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"))
}
/// Return whether a checkpoint was permanently isolated after a malformed
/// or unsupported snapshot was observed.
pub(crate) async fn is_blocked(disk: &DiskStore, task_id: &str) -> bool {
if validate_resume_task_id(task_id).is_err() {
return false;
}
let blocked_path = Self::blocked_path(task_id);
let Ok(path) = path_to_str(&blocked_path) else {
return false;
};
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
Ok(_) => true,
Err(crate::heal::DiskError::FileNotFound) => false,
Err(_) => true,
}
}
/// Validate the checkpoint while enumerating resumable state. This reads
/// the checkpoint once and also isolates malformed or unsupported data.
pub(crate) async fn is_resumable(disk: &DiskStore, task_id: &str) -> Result<bool> {
validate_resume_task_id(task_id)?;
if Self::is_blocked(disk, task_id).await {
return Err(Error::InvalidCheckpoint(format!("Resume task {task_id} has a blocked checkpoint")));
}
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
let Ok(path) = path_to_str(&file_path) else {
return Err(Error::InvalidCheckpoint("Resume checkpoint path is not valid UTF-8".to_string()));
};
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
Ok(bytes) if bytes.is_empty() => Ok(true),
Ok(bytes) => Self::load_from_data(disk.clone(), task_id, bytes.to_vec())
.await
.map(|_| true),
Err(crate::heal::DiskError::FileNotFound) => Ok(true),
Err(error) => Err(error.into()),
}
}
async fn block_invalid_snapshot(disk: &DiskStore, task_id: &str) {
// This marker is intentionally version-agnostic: an unsupported reader
// must stop selector retries until an operator cleans up the snapshot.
let blocked_path = Self::blocked_path(task_id);
let Ok(path) = path_to_str(&blocked_path) else {
return;
};
let result = EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
path,
None,
Some(EcstoreDiskBytes::from_static(b"blocked")),
)
.await;
match result {
Ok(EcstoreConditionalFileUpdate::Updated | EcstoreConditionalFileUpdate::Mismatch) => {}
Ok(EcstoreConditionalFileUpdate::Missing) => warn!(
target: "rustfs::heal::resume",
event = EVENT_HEAL_CHECKPOINT_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_RESUME,
task_id,
state = "blocked_marker_write_failed",
error = "marker target disappeared",
"Heal checkpoint could not persist its blocked marker"
),
Err(error) => warn!(
target: "rustfs::heal::resume",
event = EVENT_HEAL_CHECKPOINT_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_RESUME,
task_id,
state = "blocked_marker_write_failed",
error = %error,
"Heal checkpoint could not persist its blocked marker"
),
}
}
/// create new checkpoint manager
pub async fn new(disk: DiskStore, task_id: String) -> Result<Self> {
validate_resume_task_id(&task_id)?;
let checkpoint_volume = format!("{RUSTFS_META_BUCKET}/{BUCKET_META_PREFIX}");
if let Err(error) = EcstoreDiskAPI::make_volume(disk.as_ref(), &checkpoint_volume).await
&& error != crate::heal::DiskError::VolumeExists
{
return Err(Error::TaskExecutionFailed {
message: format!("Failed to create checkpoint volume: {error}"),
});
}
let checkpoint = ResumeCheckpoint::new(task_id);
let manager = Self {
disk,
checkpoint: Arc::new(RwLock::new(checkpoint)),
throttle: Mutex::new(PersistThrottle::new()),
save_lock: AsyncMutex::new(()),
last_saved: Mutex::new(None),
};
// save initial checkpoint
@@ -140,6 +245,7 @@ impl CheckpointManager {
error = %e,
"Heal checkpoint persistence failed"
);
return Err(e);
}
Ok(manager)
}
@@ -148,11 +254,22 @@ impl CheckpointManager {
pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result<Self> {
validate_resume_task_id(task_id)?;
let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?;
let mut checkpoint: ResumeCheckpoint =
serde_json::from_slice(&checkpoint_data).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to deserialize checkpoint: {e}"),
})?;
Self::load_from_data(disk, task_id, checkpoint_data).await
}
async fn load_from_data(disk: DiskStore, task_id: &str, checkpoint_data: Vec<u8>) -> Result<Self> {
validate_resume_task_id(task_id)?;
let mut checkpoint: ResumeCheckpoint = match serde_json::from_slice(&checkpoint_data) {
Ok(checkpoint) => checkpoint,
Err(error) => {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!("Failed to deserialize checkpoint: {error}"),
});
}
};
if checkpoint.task_id != task_id {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::TaskExecutionFailed {
message: "Resume checkpoint task id does not match filename".to_string(),
});
@@ -163,6 +280,7 @@ impl CheckpointManager {
// identities. Discard the stale sets and position, then stamp the
// current schema so the scan restarts cleanly.
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!(
"Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
@@ -170,7 +288,45 @@ impl CheckpointManager {
),
});
}
if checkpoint.schema_version < CURRENT_CHECKPOINT_SCHEMA {
let integrity_verified = if let Some(expected) = checkpoint.integrity_digest.as_deref() {
let actual = Self::checkpoint_digest(&Self::serialize_without_digest(&checkpoint)?);
if expected != actual {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::InvalidCheckpoint(format!(
"Resume checkpoint digest does not match task {task_id}"
)));
}
true
} else if checkpoint.schema_version >= CURRENT_CHECKPOINT_SCHEMA {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::InvalidCheckpoint(format!(
"Resume checkpoint digest is missing for task {task_id}"
)));
} else {
let digest_path = Self::digest_path(task_id);
let digest_path = path_to_str(&digest_path)?;
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, digest_path).await {
Ok(expected) => {
let actual = Self::checkpoint_digest(&checkpoint_data);
if expected.as_ref() != actual.as_bytes() {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::InvalidCheckpoint(format!(
"Resume checkpoint digest does not match task {task_id}"
)));
}
true
}
Err(crate::heal::DiskError::FileNotFound) => false,
Err(error) => {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to read checkpoint digest: {error}"),
});
}
}
};
if checkpoint.schema_version < CHECKPOINT_PER_VERSION_SCHEMA || !integrity_verified {
warn!(
target: "rustfs::heal::resume",
event = EVENT_HEAL_CHECKPOINT_STATE,
@@ -187,13 +343,15 @@ impl CheckpointManager {
checkpoint.skipped_objects.clear();
checkpoint.current_bucket_index = 0;
checkpoint.current_object_index = 0;
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
}
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
Ok(Self {
disk,
checkpoint: Arc::new(RwLock::new(checkpoint)),
throttle: Mutex::new(PersistThrottle::new()),
save_lock: AsyncMutex::new(()),
last_saved: Mutex::new(Some(EcstoreDiskBytes::from(checkpoint_data))),
})
}
@@ -204,7 +362,7 @@ impl CheckpointManager {
}
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
match path_to_str(&file_path) {
Ok(path_str) => match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
Ok(path_str) => match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
Ok(data) => !data.is_empty(),
Err(_) => false,
},
@@ -292,6 +450,8 @@ impl CheckpointManager {
let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
delete_resume_file(&self.disk, &checkpoint_file).await?;
delete_resume_file(&self.disk, &Self::digest_path(&task_id)).await?;
delete_resume_file(&self.disk, &Self::blocked_path(&task_id)).await?;
debug!(
target: "rustfs::heal::resume",
@@ -307,21 +467,130 @@ impl CheckpointManager {
/// save checkpoint to disk
async fn save_checkpoint(&self) -> Result<()> {
let checkpoint = self.checkpoint.read().await;
// Serialize saves and take the snapshot only after acquiring the lock:
// a slower writer must not publish a snapshot taken before a newer one.
let _save_guard = self.save_lock.lock().await;
let checkpoint = self.checkpoint.read().await.clone();
validate_resume_task_id(&checkpoint.task_id)?;
let checkpoint_data = serde_json::to_vec(&*checkpoint).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})?;
let unsigned_checkpoint_data = Self::serialize_without_digest(&checkpoint)?;
let digest = Self::checkpoint_digest(&unsigned_checkpoint_data);
let mut persisted_checkpoint = checkpoint.clone();
persisted_checkpoint.integrity_digest = Some(digest);
let checkpoint_data =
EcstoreDiskBytes::from(serde_json::to_vec(&persisted_checkpoint).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})?);
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE));
let path_str = path_to_str(&file_path)?;
self.disk
.write_all(RUSTFS_META_BUCKET, path_str, checkpoint_data.into())
let last_saved = self
.last_saved
.lock()
.map_err(|_| Error::TaskExecutionFailed {
message: "Checkpoint save state lock is poisoned; refusing to save".to_string(),
})?
.clone();
let update = EcstoreDiskAPI::compare_and_update_file(
self.disk.as_ref(),
RUSTFS_META_BUCKET,
path_str,
last_saved.clone(),
Some(checkpoint_data.clone()),
)
.await
.map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to save checkpoint: {e}"),
})?;
let expected = match update {
EcstoreConditionalFileUpdate::Updated => None,
EcstoreConditionalFileUpdate::Missing => {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
});
}
EcstoreConditionalFileUpdate::Mismatch => {
// A healthy manager normally completes the CAS above without
// another read or JSON parse. Inspect only after a mismatch so
// corruption and future schemas cannot be overwritten blindly.
let existing = match HealDiskExt::read_all(self.disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
Ok(existing) => existing,
Err(crate::heal::DiskError::FileNotFound) => {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
});
}
Err(error) => {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to inspect checkpoint after CAS mismatch: {error}"),
});
}
};
if existing.is_empty() && last_saved.is_none() {
Some(existing)
} else {
let current: ResumeCheckpoint = match serde_json::from_slice(&existing) {
Ok(current) => current,
Err(error) => {
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!("Existing checkpoint is corrupt: {error}"),
});
}
};
if current.task_id != checkpoint.task_id {
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
return Err(Error::TaskExecutionFailed {
message: "Existing checkpoint task id does not match filename".to_string(),
});
}
if current.schema_version > CURRENT_CHECKPOINT_SCHEMA {
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!(
"Existing checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
current.schema_version
),
});
}
if last_saved.as_ref().is_none_or(|saved| saved.as_ref() != existing.as_ref()) {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint changed since this manager loaded it; refusing to overwrite newer progress"
.to_string(),
});
}
Some(existing)
}
}
};
if let Some(expected) = expected {
match EcstoreDiskAPI::compare_and_update_file(
self.disk.as_ref(),
RUSTFS_META_BUCKET,
path_str,
Some(expected),
Some(checkpoint_data.clone()),
)
.await
.map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to save checkpoint: {e}"),
})?;
message: format!("Failed to save checkpoint after CAS mismatch: {e}"),
})? {
EcstoreConditionalFileUpdate::Updated => {}
EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint changed while saving; refusing to overwrite newer progress".to_string(),
});
}
}
}
let mut last_saved = self.last_saved.lock().map_err(|_| Error::TaskExecutionFailed {
message: "Checkpoint save state lock is poisoned after save".to_string(),
})?;
*last_saved = Some(checkpoint_data);
debug!(
target: "rustfs::heal::resume",
@@ -341,11 +610,38 @@ impl CheckpointManager {
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
let path_str = path_to_str(&file_path)?;
disk.read_all(RUSTFS_META_BUCKET, path_str)
HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str)
.await
.map(|bytes| bytes.to_vec())
.map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to read checkpoint file: {e}"),
})
}
fn serialize_without_digest(checkpoint: &ResumeCheckpoint) -> Result<Vec<u8>> {
let mut unsigned = checkpoint.clone();
unsigned.integrity_digest = None;
let mut value = serde_json::to_value(&unsigned).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})?;
for field in ["processed_objects", "failed_objects", "skipped_objects"] {
let Some(values) = value.get_mut(field).and_then(serde_json::Value::as_array_mut) else {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to canonicalize checkpoint field: {field}"),
});
};
values.sort_by(|left, right| left.as_str().cmp(&right.as_str()));
}
serde_json::to_vec(&value).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})
}
fn checkpoint_digest(checkpoint_data: &[u8]) -> String {
base64::engine::general_purpose::STANDARD.encode(Sha256::digest(checkpoint_data))
}
fn digest_path(task_id: &str) -> std::path::PathBuf {
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_DIGEST_FILE}"))
}
}
+389
View File
@@ -1600,6 +1600,32 @@ async fn test_checkpoint_schema_v4_discarded_on_load() {
temp_dir.close().expect("remove schema test directory");
}
#[tokio::test]
async fn downgraded_unsigned_checkpoint_resets_untrusted_progress() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
manager.add_processed_object("victim-a".to_string()).await.unwrap();
manager.update_position(2, 500).await.unwrap();
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let bytes = disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path).await.unwrap();
let mut downgraded: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
downgraded["schema_version"] = serde_json::json!(CURRENT_CHECKPOINT_SCHEMA - 1);
downgraded.as_object_mut().unwrap().remove("integrity_digest");
downgraded["processed_objects"] = serde_json::json!(["victim-b"]);
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&downgraded).unwrap().into())
.await
.expect("write downgraded checkpoint");
let manager = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
let checkpoint = manager.get_checkpoint().await;
assert_eq!(checkpoint.schema_version, CURRENT_CHECKPOINT_SCHEMA);
assert_eq!(checkpoint.current_bucket_index, 0);
assert_eq!(checkpoint.current_object_index, 0);
assert!(checkpoint.processed_objects.is_empty());
temp_dir.close().unwrap();
}
#[tokio::test]
async fn current_normal_resume_schema_preserves_progress() {
let (temp_dir, disk) = schema_test_disk().await;
@@ -1675,6 +1701,369 @@ async fn future_resume_and_checkpoint_schemas_are_rejected() {
temp_dir.close().expect("remove schema test directory");
}
#[tokio::test]
async fn checkpoint_save_does_not_replace_a_non_empty_truncated_snapshot() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let truncated = b"{\"schema_version\":5,\"task_id\":";
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, truncated.as_slice().into())
.await
.expect("write truncated checkpoint fixture");
let error = manager
.update_position(2, 7)
.await
.expect_err("a truncated checkpoint must fail closed during save");
assert!(error.to_string().contains("Existing checkpoint is corrupt"));
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read truncated checkpoint fixture"),
truncated.as_slice()
);
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove checkpoint save test directory");
}
#[tokio::test]
async fn checkpoint_save_does_not_replace_a_future_schema_snapshot() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let mut future = ResumeCheckpoint::new(task_id.clone());
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, future_bytes.clone().into())
.await
.expect("write future checkpoint fixture");
let error = manager
.update_position(2, 7)
.await
.expect_err("a future schema must fail closed during save");
assert!(error.to_string().contains("Existing checkpoint schema"));
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read future checkpoint fixture"),
future_bytes
);
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove future schema test directory");
}
#[tokio::test]
async fn checkpoint_digest_rejects_same_length_progress_tampering() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
manager
.add_processed_object("victim-a".to_string())
.await
.expect("persist checkpoint progress");
manager.update_position(1, 1).await.expect("flush checkpoint progress");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let original = disk
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read checkpoint fixture");
let tampered = original
.windows(b"victim-a".len())
.position(|window| window == b"victim-a")
.map(|index| {
let mut bytes = original.to_vec();
bytes[index..index + b"victim-a".len()].copy_from_slice(b"victim-b");
bytes
})
.expect("checkpoint should contain the processed object");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, tampered.into())
.await
.expect("write tampered checkpoint fixture");
assert!(CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err());
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove digest test directory");
}
#[tokio::test]
async fn checkpoint_integrity_survives_missing_legacy_sidecar() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
manager.update_position(2, 9).await.unwrap();
let digest_path = format!("{BUCKET_META_PREFIX}/{task_id}_ahm_checkpoint.sha256");
delete_resume_file(&disk, Path::new(&digest_path)).await.unwrap();
let restored = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
let checkpoint = restored.get_checkpoint().await;
assert_eq!(checkpoint.current_bucket_index, 2);
assert_eq!(checkpoint.current_object_index, 9);
assert!(checkpoint.integrity_digest.is_some());
temp_dir.close().unwrap();
}
#[tokio::test]
async fn checkpoint_integrity_survives_multi_object_reload() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
for index in 0..32 {
manager.add_processed_object(format!("processed-{index}")).await.unwrap();
manager.add_failed_object(format!("failed-{index}")).await.unwrap();
manager.add_skipped_object(format!("skipped-{index}")).await.unwrap();
}
manager.update_position(2, 9).await.unwrap();
CheckpointManager::load_from_disk(disk, &task_id)
.await
.expect("a healthy multi-object checkpoint must survive reload");
temp_dir.close().unwrap();
}
#[tokio::test]
async fn checkpoint_integrity_rejects_a_removed_embedded_digest() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
manager.update_position(2, 9).await.unwrap();
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let bytes = disk
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read checkpoint fixture");
let mut value: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
value["current_object_index"] = serde_json::json!(10);
value.as_object_mut().unwrap().remove("integrity_digest");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&value).unwrap().into())
.await
.expect("write tampered checkpoint fixture");
assert!(
CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err(),
"a current checkpoint without its embedded digest must fail closed"
);
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().unwrap();
}
#[tokio::test]
async fn new_checkpoint_manager_rebuilds_an_empty_snapshot() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, EcstoreDiskBytes::new())
.await
.expect("write empty checkpoint fixture");
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("a new manager must rebuild an empty checkpoint");
manager
.update_position(3, 11)
.await
.expect("rebuilt checkpoint must remain writable");
assert!(CheckpointManager::has_checkpoint(&disk, &task_id).await);
temp_dir.close().expect("remove empty checkpoint test directory");
}
#[tokio::test]
async fn deleted_checkpoint_is_not_recreated_by_an_old_manager() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
manager.cleanup().await.expect("delete checkpoint fixture");
let error = manager
.update_position(1, 2)
.await
.expect_err("an old manager must not resurrect a deleted checkpoint");
assert!(error.to_string().contains("removed after this manager saved it"));
assert!(!CheckpointManager::has_checkpoint(&disk, &task_id).await);
temp_dir.close().expect("remove deleted checkpoint test directory");
}
#[cfg(unix)]
#[tokio::test]
async fn checkpoint_cleanup_leaves_no_task_specific_lock_artifact() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let lock_path = Path::new(BUCKET_META_PREFIX)
.join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"))
.with_extension("rustfs-cas.lock");
let lock_path = temp_dir.path().join(RUSTFS_META_BUCKET).join(lock_path);
manager.cleanup().await.expect("delete checkpoint fixture");
assert!(
!lock_path.exists(),
"successful checkpoint cleanup must not leave a task-specific lock artifact"
);
}
#[tokio::test]
async fn an_empty_blocked_marker_still_blocks_resume_selection() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &blocked_path, EcstoreDiskBytes::new())
.await
.expect("write empty blocked marker fixture");
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
assert!(CheckpointManager::is_resumable(&disk, &task_id).await.is_err());
// Recovery requires replacing/cleaning the snapshot, then removing the
// marker; ordinary selector retries are intentionally not an unlock path.
manager.cleanup().await.expect("clean blocked checkpoint");
assert!(!CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove empty blocked marker test directory");
}
#[tokio::test]
async fn resumable_selector_skips_healthy_tasks_with_blocked_markers() {
let (temp_dir, disk) = schema_test_disk().await;
let tasks = [
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::new()),
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::from_static(b"blocked")),
];
for (task_id, marker) in &tasks {
ResumeManager::new(
disk.clone(),
task_id.clone(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
vec!["bucket".to_string()],
)
.await
.expect("create healthy resume state");
CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create healthy checkpoint");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let checkpoint_bytes = disk
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read healthy checkpoint before blocking");
let marker_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &marker_path, marker.clone())
.await
.expect("write blocked marker");
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read healthy checkpoint after blocking"),
checkpoint_bytes
);
}
temp_dir.close().expect("remove blocked selector test directory");
}
#[tokio::test]
async fn stale_checkpoint_manager_cannot_overwrite_newer_progress() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let first = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create first checkpoint manager");
let second = CheckpointManager::load_from_disk(disk.clone(), &task_id)
.await
.expect("load second checkpoint manager");
second
.update_position(4, 20)
.await
.expect("persist newer checkpoint progress");
let error = first
.update_position(1, 3)
.await
.expect_err("stale checkpoint manager must not overwrite newer progress");
assert!(error.to_string().contains("newer progress"));
let persisted = CheckpointManager::load_from_disk(disk.clone(), &task_id)
.await
.expect("load newer checkpoint progress")
.get_checkpoint()
.await;
assert_eq!(persisted.current_bucket_index, 4);
assert_eq!(persisted.current_object_index, 20);
temp_dir.close().expect("remove stale manager test directory");
}
#[tokio::test]
async fn resumable_selector_isolates_future_and_corrupt_checkpoints() {
let (temp_dir, disk) = schema_test_disk().await;
let future_task = ResumeUtils::generate_task_id();
let corrupt_task = ResumeUtils::generate_task_id();
for task_id in [&future_task, &corrupt_task] {
ResumeManager::new(
disk.clone(),
task_id.to_string(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
vec!["bucket".to_string()],
)
.await
.expect("create resumable state fixture");
}
let future_path = format!("{BUCKET_META_PREFIX}/{future_task}_{RESUME_CHECKPOINT_FILE}");
let mut future = ResumeCheckpoint::new(future_task.clone());
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
disk.write_all(RUSTFS_META_BUCKET, &future_path, future_bytes.clone().into())
.await
.expect("write future checkpoint fixture");
let corrupt_path = format!("{BUCKET_META_PREFIX}/{corrupt_task}_{RESUME_CHECKPOINT_FILE}");
let corrupt_bytes = b"{truncated";
disk.write_all(RUSTFS_META_BUCKET, &corrupt_path, corrupt_bytes.as_slice().into())
.await
.expect("write corrupt checkpoint fixture");
assert!(CheckpointManager::is_resumable(&disk, &future_task).await.is_err());
assert!(CheckpointManager::is_resumable(&disk, &corrupt_task).await.is_err());
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
for (task_id, path, bytes) in [
(&future_task, future_path, future_bytes),
(&corrupt_task, corrupt_path, corrupt_bytes.to_vec()),
] {
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("read isolated checkpoint bytes"),
bytes
);
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
assert!(
!disk
.read_all(RUSTFS_META_BUCKET, &blocked_path)
.await
.expect("read checkpoint blocked marker")
.is_empty()
);
}
temp_dir.close().expect("remove selector isolation test directory");
}
#[test]
fn test_persist_throttle_batches_until_threshold() {
let mut throttle = PersistThrottle::new();
+2 -1
View File
@@ -21,7 +21,7 @@ use uuid::Uuid;
use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord};
use super::{
EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
CheckpointManager, EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str,
replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id,
};
@@ -67,6 +67,7 @@ impl ResumeUtils {
// Extract task ID from filename: {task_id}_ahm_resume_state.json
if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}"))
&& validate_resume_task_id(task_id).is_ok()
&& CheckpointManager::is_resumable(disk, task_id).await?
{
task_ids.push(task_id.to_string());
}
+25 -86
View File
@@ -16,14 +16,14 @@
//!
//! `scripts/test/vault_ha_kms_live.sh` owns the official Vault containers and
//! kills the active node while this test continuously decrypts through a
//! surviving standby. KV2 and Transit must recover after the bounded circuit
//! interval, use a bounded number of attempts, and leave the circuit and
//! in-flight gauges at zero after a new leader is elected.
//! surviving standby. KV2 and Transit requests must remain successful, use a
//! bounded number of attempts, and leave the circuit and in-flight gauges at
//! zero after a new leader is elected.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use metrics_util::MetricKind;
@@ -43,11 +43,6 @@ const OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
const IN_FLIGHT: &str = "rustfs_kms_backend_in_flight";
const CIRCUIT_OPEN: &str = "rustfs_kms_backend_circuit_open";
const MAX_ATTEMPTS: u32 = 10;
const ATTEMPT_TIMEOUT: Duration = Duration::from_secs(2);
const HEALTHY_PROGRESS_TIMEOUT: Duration = Duration::from_secs(20);
// The circuit remains open for 30s after five failed attempts.
const POST_FAILOVER_PROGRESS_TIMEOUT: Duration = Duration::from_secs(35);
const FAILOVER_ERROR_POLL_INTERVAL: Duration = Duration::from_millis(100);
type MetricEntry = (
metrics_util::CompositeKey,
@@ -69,7 +64,7 @@ fn config(backend: KmsBackend, backend_config: BackendConfig) -> KmsConfig {
backend,
backend_config,
allow_insecure_dev_defaults: true,
timeout: ATTEMPT_TIMEOUT,
timeout: Duration::from_secs(2),
retry_attempts: MAX_ATTEMPTS,
enable_cache: false,
..KmsConfig::default()
@@ -169,31 +164,14 @@ fn retryable_failures(snapshot: &[MetricEntry], operation: &str) -> u64 {
.sum()
}
async fn wait_for_count(
counter: &AtomicU64,
failure: &Mutex<Option<String>>,
minimum: u64,
description: &str,
timeout: Duration,
) {
tokio::time::timeout(timeout, async {
async fn wait_for_count(counter: &AtomicU64, minimum: u64, description: &str) {
tokio::time::timeout(Duration::from_secs(20), async {
while counter.load(Ordering::SeqCst) < minimum {
if let Some(error) = failure.lock().expect("decrypt failure lock poisoned").as_ref() {
panic!(
"{description} worker failed after {} successful decrypts: {error}",
counter.load(Ordering::SeqCst)
);
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.unwrap_or_else(|_| {
panic!(
"timed out after {timeout:?} waiting for {description}: completed {}, expected {minimum}",
counter.load(Ordering::SeqCst)
)
});
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
}
async fn wait_for_file(path: &Path, description: &str) {
@@ -211,8 +189,7 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
request: DecryptRequest,
expected: Vec<u8>,
completed: Arc<AtomicU64>,
allow_failover_errors: Arc<AtomicBool>,
failure: Arc<Mutex<Option<String>>>,
failed: Arc<AtomicBool>,
stop: CancellationToken,
) {
while !stop.is_cancelled() {
@@ -220,18 +197,8 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
Ok(response) if response.plaintext == expected => {
completed.fetch_add(1, Ordering::SeqCst);
}
Ok(_) => {
*failure.lock().expect("decrypt failure lock poisoned") =
Some("decrypt returned unexpected plaintext".to_string());
return;
}
Err(rustfs_kms::KmsError::BackendError { .. } | rustfs_kms::KmsError::OperationTimedOut { .. })
if allow_failover_errors.load(Ordering::SeqCst) =>
{
tokio::time::sleep(FAILOVER_ERROR_POLL_INTERVAL).await;
}
Err(error) => {
*failure.lock().expect("decrypt failure lock poisoned") = Some(error.to_string());
Ok(_) | Err(_) => {
failed.store(true, Ordering::SeqCst);
return;
}
}
@@ -329,9 +296,7 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
);
let stop = CancellationToken::new();
let allow_failover_errors = Arc::new(AtomicBool::new(false));
let kv2_failure = Arc::new(Mutex::new(None));
let transit_failure = Arc::new(Mutex::new(None));
let failed = Arc::new(AtomicBool::new(false));
let kv2_completed = Arc::new(AtomicU64::new(0));
let transit_completed = Arc::new(AtomicU64::new(0));
let kv2_worker = tokio::spawn(decrypt_loop(
@@ -339,8 +304,7 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
kv2_request,
kv2_data_key.plaintext_key,
Arc::clone(&kv2_completed),
Arc::clone(&allow_failover_errors),
Arc::clone(&kv2_failure),
Arc::clone(&failed),
stop.clone(),
));
let transit_worker = tokio::spawn(decrypt_loop(
@@ -348,21 +312,12 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
transit_request,
transit_data_key.plaintext_key,
Arc::clone(&transit_completed),
Arc::clone(&allow_failover_errors),
Arc::clone(&transit_failure),
Arc::clone(&failed),
stop.clone(),
));
wait_for_count(&kv2_completed, &kv2_failure, 2, "two healthy KV2 decrypts", HEALTHY_PROGRESS_TIMEOUT).await;
wait_for_count(
&transit_completed,
&transit_failure,
2,
"two healthy Transit decrypts",
HEALTHY_PROGRESS_TIMEOUT,
)
.await;
allow_failover_errors.store(true, Ordering::SeqCst);
wait_for_count(&kv2_completed, 2, "two healthy KV2 decrypts").await;
wait_for_count(&transit_completed, 2, "two healthy Transit decrypts").await;
std::fs::write(&marker, b"ready").expect("publish failover readiness marker");
wait_for_file(&elected, "the replacement Vault leader").await;
@@ -371,39 +326,18 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
let kv2_after_election = kv2_completed.load(Ordering::SeqCst) + 2;
let transit_after_election = transit_completed.load(Ordering::SeqCst) + 2;
wait_for_count(
&kv2_completed,
&kv2_failure,
kv2_after_election,
"post-failover KV2 decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(
&transit_completed,
&transit_failure,
transit_after_election,
"post-failover Transit decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(&kv2_completed, kv2_after_election, "post-failover KV2 decrypts").await;
wait_for_count(&transit_completed, transit_after_election, "post-failover Transit decrypts").await;
stop.cancel();
kv2_worker.await.expect("KV2 decrypt worker must join");
transit_worker.await.expect("Transit decrypt worker must join");
assert!(
kv2_failure.lock().expect("KV2 failure lock poisoned").is_none(),
"no KV2 decrypt may fail or return different plaintext"
);
assert!(
transit_failure.lock().expect("Transit failure lock poisoned").is_none(),
"no Transit decrypt may fail or return different plaintext"
);
assert!(!failed.load(Ordering::SeqCst), "no decrypt may fail or return different plaintext");
}
#[test]
#[ignore = "requires a real three-node Vault Raft cluster; run scripts/test/vault_ha_kms_live.sh"]
fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
@@ -415,6 +349,11 @@ fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
});
let snapshot = snapshotter.snapshot().into_vec();
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "circuit_open")]),
0,
"a bounded leader election must not open the circuit"
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]),
0,
@@ -153,99 +153,6 @@ No migration step is required for these decisions because this note documents th
current RustFS behavior. Changing either decision later requires an operator
compatibility note and updated characterization tests.
## Tier Free Versions During Decommission
A tier free version is an internal xl.meta record (`rustfs_filemeta::FREE_VERSION`,
flagged `XL_FLAG_FREE_VERSION`) shaped like a delete marker. It is created by
`MetaObject::init_free_version` when a version whose remote transition completed is
deleted locally: the visible version is removed and the record keeps the remote-tier
identity (tier, object name, version id, state, destination id) needed for an
idempotent remote delete. Free versions are not user-visible versions; `num_versions`
and all listing/GET paths exclude them.
### Lifecycle And Consumers
Creation: any local delete that removes a version whose transition status is
`complete` appends the record via `MetaObject::delete_version`
`init_free_version` (skipped only when `skip_tier_free_version` is set, as on
data-movement copies). The same deletes also persist a durable tier-journal
entry on every user-facing path: S3 single deletes (`execute_delete_object`
`delete_object_with_tier_delete_journal`), S3 batch deletes, lifecycle expiry,
and lifecycle delete-all all prepare and commit a journal entry around the
delete. A journal entry is omitted when the removed version's transition state
decodes as `TransitionVersionState::Unknown`, or on internal journal-less
delete paths that never touch transitioned user objects.
Consumption while the record exists: the background recovery loop started by
`init_background_expiry` (spawned by `spawn_tier_free_version_recovery_once`,
enabled by default) scans disks for pending records and re-enqueues them; the
usage scanner does the same; the lifecycle worker then deletes the remote tier
object idempotently and only afterwards removes the local record. Heal walks
include free-version records in metadata healing. Transition planning,
replication, restore, GET, listings, and usage aggregation never depend on
them.
### Decommission Handling
The exact decommission inventory loader (`load_file_info_versions_exact` via
`get_all_file_info_versions`) keeps free-version records inline in `versions`.
The migration loop handles them before lifecycle expiry and delete-marker
shortcuts. It selects a target pool using the free-version-aware lookup, then
writes the original free record to every target disk with the normal metadata
write quorum. The free-version marker, local version id, transition identity,
transition state, and destination id are preserved at the FileInfo/metadata
boundary.
The source record is physically removed only after the target write quorum has
committed and the source cleanup preflight still matches the exact inventory.
If the lifecycle worker has already completed the remote delete and removed the
source record before decommission acquires the source lock, decommission records
that identity as already consumed and treats the missing source record as safe.
If target capacity, metadata validation, lock fencing, or quorum fails, the
source record remains and the entry records `state = "free_version_retained"`
with reason `tier_free_version_migration_failed`; the worker retries the
operation on a later pass. A target record with the same version id is accepted
only when its free-version identity matches; a conflicting ordinary version or
different free record is an overwrite error. This makes retries idempotent and
prevents a free record from replacing a user-visible version.
### Reference-Audit Result
After migration, user-facing GET/list/transition/replication/restore paths still
exclude the record. Recovery, usage scanning, lifecycle tier cleanup, and heal
continue to see it when they request free versions, so an unresolved remote
delete remains actionable on the target pool. The committed tier journal remains
an independent retry source where one exists; it is not used as a reason to drop
the xl.meta record. In particular, `Unknown` transition state records are
migrated unchanged rather than discarded: the lifecycle worker retains them if
remote identity validation cannot make a delete request.
Each migrated record emits `state = "free_version_migrated"` with reason
`tier_free_version_migrated`. A record consumed before migration emits
`state = "free_version_consumed"` with reason
`tier_free_version_already_consumed`. Each failed record emits the retained state
and failure reason above. The entry also emits a disposition summary with
migrated, consumed, retained, and total counts. The final decommission sweep uses
the exact loader, counts free records still present, and emits one retained
record/reason for each unresolved free version before failing the sweep. This
makes successful migration, completed cleanup, and retained cleanup obligations
visible instead of silently omitting free records.
No new S3-visible version or admin response field is needed: free versions remain
internal and are never counted as user-visible versions. The structured
`decommission_entry` events are the operational status surface for the
free-version disposition; the existing decommission item/failed counters still
report the enclosing object migration result.
Regression guard:
- `decommission_tier_free_version_preserves_remote_identity`
- `decommission_tier_free_version_resume_requires_write_quorum`
- `decommission_tier_free_version_commit_rejects_lost_fence`
- `test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source`
- `decommission_entry_skips_cleanup_only_marker_when_free_version_is_present`
- `decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source`
## Regression Guard
The queued multi-pool contract is guarded by:
+1 -1
View File
@@ -241,7 +241,7 @@ env \
RUSTFS_TEST_VAULT_FAILOVER_MARKER="$MARKER" \
RUSTFS_TEST_VAULT_OLD_LEADER="$OLD_LEADER" \
cargo test -p rustfs-kms --test vault_ha_failover_live \
vault_raft_leader_failure_recovers_kv2_and_transit_decrypts -- \
vault_raft_leader_failure_preserves_kv2_and_transit_decrypts -- \
--ignored --nocapture --test-threads=1 &
TEST_PID=$!