mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 20:43:04 +00:00
refactor(storage): bind rename rollback to its commit context (#8324)
This commit is contained in:
@@ -84,46 +84,107 @@ async fn wait_for_heal_commit_test_barrier(root: &Path, bucket: &str, object: &s
|
||||
}
|
||||
}
|
||||
|
||||
fn rollback_committed_rename_std(
|
||||
dst_file_path: &Path,
|
||||
new_data_path: Option<&Path>,
|
||||
rollback_data_dir: Option<Uuid>,
|
||||
) -> std::io::Result<()> {
|
||||
if let Some(old_data_dir) = rollback_data_dir {
|
||||
let Some(dst_parent) = dst_file_path.parent() else {
|
||||
return Err(std::io::Error::new(ErrorKind::InvalidInput, "missing object metadata parent"));
|
||||
};
|
||||
let backup_path = dst_parent.join(old_data_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP);
|
||||
std::fs::rename(backup_path, dst_file_path)?;
|
||||
} else {
|
||||
remove_file_if_exists(dst_file_path)?;
|
||||
}
|
||||
|
||||
if let Some(new_data_path) = new_data_path {
|
||||
remove_dir_all_if_exists(new_data_path)?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
struct RenameCommitContext<'a> {
|
||||
// Drop the directory identity guard before releasing mutation ownership.
|
||||
directory_guard: os::RenameCommitGuard,
|
||||
publication_root: &'a os::PublicationRoot,
|
||||
mutation_lease: Arc<os::NamespaceMutationLease>,
|
||||
source_metadata: &'a Path,
|
||||
destination_metadata: &'a Path,
|
||||
source_volume_dir: &'a Path,
|
||||
data_paths: Option<&'a (PathBuf, PathBuf)>,
|
||||
}
|
||||
|
||||
fn rollback_inline_metadata_commit_std(
|
||||
dst_file_path: &Path,
|
||||
rollback_data_dir: Option<Uuid>,
|
||||
local_rollback_path: Option<&Path>,
|
||||
) -> std::io::Result<()> {
|
||||
#[cfg(all(test, not(windows)))]
|
||||
os::prepared_publication_test_hooks::run(os::prepared_publication_test_hooks::Stage::Rollback, dst_file_path);
|
||||
if let Some(backup_path) = local_rollback_path {
|
||||
// The commit immediately before this rollback renamed the staged
|
||||
// xl.meta from the same directory as `backup_path` onto
|
||||
// `dst_file_path`, proving both paths are on the same filesystem.
|
||||
// Unix rename atomically replaces the committed destination; never
|
||||
// unlink it first or an interrupted rollback could lose xl.meta.
|
||||
std::fs::rename(backup_path, dst_file_path)?;
|
||||
} else {
|
||||
rollback_committed_rename_std(dst_file_path, None, rollback_data_dir)?;
|
||||
struct StagedRollbackBackup(PathBuf);
|
||||
|
||||
impl StagedRollbackBackup {
|
||||
fn path(&self) -> &Path {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl RenameCommitContext<'_> {
|
||||
fn rollback_metadata_std(&self, rollback_data_dir: Option<Uuid>) -> std::io::Result<()> {
|
||||
if let Some(old_data_dir) = rollback_data_dir {
|
||||
let Some(dst_parent) = self.destination_metadata.parent() else {
|
||||
return Err(std::io::Error::new(ErrorKind::InvalidInput, "missing object metadata parent"));
|
||||
};
|
||||
let backup_path = dst_parent.join(old_data_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP);
|
||||
std::fs::rename(backup_path, self.destination_metadata)?;
|
||||
} else {
|
||||
remove_file_if_exists(self.destination_metadata)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn rollback_non_inline_std(&self, rollback_data_dir: Option<Uuid>) -> std::io::Result<()> {
|
||||
self.rollback_metadata_std(rollback_data_dir)?;
|
||||
if let Some((_, new_data_path)) = self.data_paths {
|
||||
remove_dir_all_if_exists(new_data_path)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn rollback_inline_std(
|
||||
&self,
|
||||
rollback_data_dir: Option<Uuid>,
|
||||
local_backup: Option<&StagedRollbackBackup>,
|
||||
) -> std::io::Result<()> {
|
||||
#[cfg(all(test, not(windows)))]
|
||||
os::prepared_publication_test_hooks::run(os::prepared_publication_test_hooks::Stage::Rollback, self.destination_metadata);
|
||||
if let Some(backup) = local_backup {
|
||||
// The preceding commit used this staged directory on the same
|
||||
// filesystem. Never unlink xl.meta before its atomic replacement.
|
||||
std::fs::rename(backup.path(), self.destination_metadata)?;
|
||||
} else {
|
||||
self.rollback_metadata_std(rollback_data_dir)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn restore_published_data_source(&self) -> Result<()> {
|
||||
let Some((src_data_path, dst_data_path)) = self.data_paths else {
|
||||
return Ok(());
|
||||
};
|
||||
let src_volume_dir = self.source_volume_dir;
|
||||
let publication_root = self.publication_root;
|
||||
let mutation_lease = self.mutation_lease.clone();
|
||||
if fs::symlink_metadata(src_data_path).await.is_ok() {
|
||||
return Ok(());
|
||||
}
|
||||
let result =
|
||||
match os::rename_all_with_lease(dst_data_path, src_data_path, src_volume_dir, publication_root, mutation_lease).await
|
||||
{
|
||||
Ok(()) => Ok(()),
|
||||
Err(DiskError::FileNotFound) => {
|
||||
let source_exists = fs::symlink_metadata(src_data_path).await.is_ok();
|
||||
let destination_missing = matches!(
|
||||
fs::symlink_metadata(dst_data_path).await,
|
||||
Err(err) if err.kind() == ErrorKind::NotFound
|
||||
);
|
||||
if source_exists && destination_missing {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(DiskError::FileNotFound)
|
||||
}
|
||||
}
|
||||
Err(err) => Err(err),
|
||||
};
|
||||
if let Err(err) = &result {
|
||||
warn!(
|
||||
target: "rustfs_ecstore::disk::local",
|
||||
event = EVENT_DISK_LOCAL_RENAME_REJECTED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_DISK_LOCAL,
|
||||
reason = "restore_staged_data_source_failed",
|
||||
src_path = ?src_data_path,
|
||||
dst_path = ?dst_data_path,
|
||||
error = ?err,
|
||||
"Failed to restore staged data after a metadata commit was rejected"
|
||||
);
|
||||
}
|
||||
result
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(any(not(windows), test))]
|
||||
@@ -212,61 +273,6 @@ async fn read_rename_destination_metadata(
|
||||
.map_err(DiskError::from)
|
||||
}
|
||||
|
||||
async fn restore_renamed_data_source(
|
||||
src_volume_dir: &Path,
|
||||
src_data_path: &Path,
|
||||
dst_data_path: &Path,
|
||||
publication_root: &os::PublicationRoot,
|
||||
mutation_lease: Arc<os::NamespaceMutationLease>,
|
||||
) -> Result<()> {
|
||||
if fs::symlink_metadata(src_data_path).await.is_ok() {
|
||||
return Ok(());
|
||||
}
|
||||
let result =
|
||||
match os::rename_all_with_lease(dst_data_path, src_data_path, src_volume_dir, publication_root, mutation_lease).await {
|
||||
Ok(()) => Ok(()),
|
||||
Err(DiskError::FileNotFound) => {
|
||||
let source_exists = fs::symlink_metadata(src_data_path).await.is_ok();
|
||||
let destination_missing = matches!(
|
||||
fs::symlink_metadata(dst_data_path).await,
|
||||
Err(err) if err.kind() == ErrorKind::NotFound
|
||||
);
|
||||
if source_exists && destination_missing {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(DiskError::FileNotFound)
|
||||
}
|
||||
}
|
||||
Err(err) => Err(err),
|
||||
};
|
||||
if let Err(err) = &result {
|
||||
warn!(
|
||||
target: "rustfs_ecstore::disk::local",
|
||||
event = EVENT_DISK_LOCAL_RENAME_REJECTED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_DISK_LOCAL,
|
||||
reason = "restore_staged_data_source_failed",
|
||||
src_path = ?src_data_path,
|
||||
dst_path = ?dst_data_path,
|
||||
error = ?err,
|
||||
"Failed to restore staged data after a metadata commit was rejected"
|
||||
);
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
async fn restore_published_data_source(
|
||||
data_paths: Option<&(PathBuf, PathBuf)>,
|
||||
src_volume_dir: &Path,
|
||||
publication_root: &os::PublicationRoot,
|
||||
mutation_lease: Arc<os::NamespaceMutationLease>,
|
||||
) -> Result<()> {
|
||||
let Some((src_data_path, dst_data_path)) = data_paths else {
|
||||
return Ok(());
|
||||
};
|
||||
restore_renamed_data_source(src_volume_dir, src_data_path, dst_data_path, publication_root, mutation_lease).await
|
||||
}
|
||||
|
||||
/// Proof produced only when the local rename returns at an existing access
|
||||
/// preflight, before metadata, backups, or object data can be published.
|
||||
#[derive(Debug)]
|
||||
@@ -457,7 +463,18 @@ impl LocalDisk {
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
let has_dst_buf = read_rename_destination_metadata(&dst_file_path, &rename_commit_guard, mutation_lease.clone()).await?;
|
||||
let mut commit = RenameCommitContext {
|
||||
directory_guard: rename_commit_guard,
|
||||
publication_root: &self.publication_root,
|
||||
mutation_lease,
|
||||
source_metadata: &src_file_path,
|
||||
destination_metadata: &dst_file_path,
|
||||
source_volume_dir: &src_volume_dir,
|
||||
data_paths: has_data_dir_path.as_ref(),
|
||||
};
|
||||
let has_dst_buf =
|
||||
read_rename_destination_metadata(commit.destination_metadata, &commit.directory_guard, commit.mutation_lease.clone())
|
||||
.await?;
|
||||
|
||||
if no_inline {
|
||||
// Non-inline: read xl.meta, parse, write, rename data dir, rename xl.meta
|
||||
@@ -533,18 +550,18 @@ impl LocalDisk {
|
||||
// rename below to report through the existing rollback path. Payload
|
||||
// durability is kept by both strict and relaxed.
|
||||
let tmp_meta_write = {
|
||||
let src_file_path = src_file_path.clone();
|
||||
let src_file_path = commit.source_metadata.to_path_buf();
|
||||
let dst_file_path = dst_file_path.clone();
|
||||
let rename_commit_guard = rename_commit_guard.clone();
|
||||
let mutation_lease = mutation_lease.clone();
|
||||
let owned_directory_guard = commit.directory_guard.clone();
|
||||
let owned_mutation_lease = commit.mutation_lease.clone();
|
||||
async move {
|
||||
os::run_blocking_namespace_operation(mutation_lease, move || {
|
||||
os::run_blocking_namespace_operation(owned_mutation_lease, move || {
|
||||
#[cfg(test)]
|
||||
run_owned_file_write_before_open(&src_file_path);
|
||||
let mut prepared_metadata_source = os::create_prepared_rename_source_with_commit_guard(
|
||||
&src_file_path,
|
||||
&dst_file_path,
|
||||
&rename_commit_guard,
|
||||
&owned_directory_guard,
|
||||
)?;
|
||||
prepared_metadata_source.write_all(&new_dst_buf, tmp_meta_sync != SyncMode::None)?;
|
||||
Ok(prepared_metadata_source)
|
||||
@@ -570,14 +587,14 @@ impl LocalDisk {
|
||||
// sequential version did.
|
||||
let prepared_metadata_source = tmp_meta_res?;
|
||||
shard_sync_res?;
|
||||
let rename_commit_guard = remove_dst_base_before_commit(
|
||||
commit.directory_guard = remove_dst_base_before_commit(
|
||||
dst_path,
|
||||
rename_commit_guard,
|
||||
commit.directory_guard,
|
||||
src_file_parent,
|
||||
dst_file_parent,
|
||||
&dst_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
commit.publication_root,
|
||||
commit.mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
if should_remove_staged_meta_before_commit(dst_path) {
|
||||
@@ -595,7 +612,7 @@ impl LocalDisk {
|
||||
if fi_healing
|
||||
&& let Some((_, dst_data_path)) = has_data_dir_path.as_ref()
|
||||
&& let Err(err) = self
|
||||
.move_to_trash_with_namespace_owner(dst_data_path, true, false, Some(mutation_lease.clone()))
|
||||
.move_to_trash_with_namespace_owner(dst_data_path, true, false, Some(commit.mutation_lease.clone()))
|
||||
.await
|
||||
{
|
||||
warn!(
|
||||
@@ -613,9 +630,9 @@ impl LocalDisk {
|
||||
src_data_path,
|
||||
dst_data_path,
|
||||
&skip_parent,
|
||||
&self.publication_root,
|
||||
&rename_commit_guard,
|
||||
mutation_lease.clone(),
|
||||
commit.publication_root,
|
||||
&commit.directory_guard,
|
||||
commit.mutation_lease.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -630,13 +647,7 @@ impl LocalDisk {
|
||||
error = ?err,
|
||||
"Disk local rename flow failed"
|
||||
);
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(err);
|
||||
}
|
||||
#[cfg(test)]
|
||||
@@ -661,13 +672,7 @@ impl LocalDisk {
|
||||
reason = "test_fail_before_old_metadata_backup",
|
||||
"Disk local rename flow failed before metadata commit"
|
||||
);
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(DiskError::Unexpected);
|
||||
}
|
||||
|
||||
@@ -685,25 +690,16 @@ impl LocalDisk {
|
||||
let backup_parent = dst_file_parent.join(old_data_dir.to_string());
|
||||
#[cfg(not(windows))]
|
||||
if let Err(err) = os::make_dir_all(&backup_parent, &skip_parent).await {
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(err);
|
||||
}
|
||||
let backup_path_guard = match rename_commit_guard.create_destination_directory_for_path_access(&backup_parent) {
|
||||
let backup_path_guard = match commit
|
||||
.directory_guard
|
||||
.create_destination_directory_for_path_access(&backup_parent)
|
||||
{
|
||||
Ok(guard) => guard,
|
||||
Err(err) => {
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(DiskError::from(to_file_error(err)));
|
||||
}
|
||||
};
|
||||
@@ -711,13 +707,7 @@ impl LocalDisk {
|
||||
if let Err(err) = check_path_length(backup_path.to_string_lossy().as_ref()) {
|
||||
#[cfg(windows)]
|
||||
drop(backup_path_guard);
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(err);
|
||||
}
|
||||
let backup_bytes = dst_buf.clone();
|
||||
@@ -726,7 +716,7 @@ impl LocalDisk {
|
||||
// sync finish. A detached spawn_blocking writer could survive
|
||||
// cancellation and later truncate a newer transaction's
|
||||
// deterministic rollback backup.
|
||||
let write_result = os::run_blocking_namespace_operation(mutation_lease.clone(), move || {
|
||||
let write_result = os::run_blocking_namespace_operation(commit.mutation_lease.clone(), move || {
|
||||
#[cfg(test)]
|
||||
run_owned_file_write_before_open(&backup_path);
|
||||
backup_path_guard.write_file_for_path_access(
|
||||
@@ -749,13 +739,7 @@ impl LocalDisk {
|
||||
error = ?err,
|
||||
"Disk local rename flow failed"
|
||||
);
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(err);
|
||||
}
|
||||
}
|
||||
@@ -773,9 +757,9 @@ impl LocalDisk {
|
||||
&src_file_path,
|
||||
&dst_file_path,
|
||||
&skip_parent,
|
||||
&self.publication_root,
|
||||
&rename_commit_guard,
|
||||
mutation_lease.clone(),
|
||||
commit.publication_root,
|
||||
&commit.directory_guard,
|
||||
commit.mutation_lease.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -790,20 +774,11 @@ impl LocalDisk {
|
||||
error = ?err,
|
||||
"Disk local rename flow failed"
|
||||
);
|
||||
restore_published_data_source(
|
||||
has_data_dir_path.as_ref(),
|
||||
&src_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
commit.restore_published_data_source().await?;
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
let committed_new_data_path = has_data_dir_path.as_ref().map(|(_, dst_data_path)| dst_data_path.as_path());
|
||||
if should_fail_after_metadata_commit(dst_path) {
|
||||
rollback_committed_rename_std(&dst_file_path, committed_new_data_path, rollback_data_dir)
|
||||
.map_err(to_file_error)?;
|
||||
commit.rollback_non_inline_std(rollback_data_dir).map_err(to_file_error)?;
|
||||
return Err(DiskError::Unexpected);
|
||||
}
|
||||
|
||||
@@ -823,13 +798,12 @@ impl LocalDisk {
|
||||
&& let Some(parent) = dst_file_path.parent()
|
||||
{
|
||||
let fsync_started = rustfs_io_metrics::put_stage_timer();
|
||||
if let Err(err) = os::fsync_dst_dir_group_commit(parent, Some(mutation_lease.clone())).await {
|
||||
if let Err(err) = os::fsync_dst_dir_group_commit(parent, Some(commit.mutation_lease.clone())).await {
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_DST_DIR_FSYNC,
|
||||
fsync_started,
|
||||
);
|
||||
rollback_committed_rename_std(&dst_file_path, committed_new_data_path, rollback_data_dir)
|
||||
.map_err(to_file_error)?;
|
||||
commit.rollback_non_inline_std(rollback_data_dir).map_err(to_file_error)?;
|
||||
// The commit rename changed the dst part inodes before this fsync
|
||||
// failed and rolled them back; drop any fd cached during that
|
||||
// window so readers re-open the restored inode (rustfs/backlog#1177).
|
||||
@@ -861,13 +835,12 @@ impl LocalDisk {
|
||||
break;
|
||||
}
|
||||
let fsync_started = rustfs_io_metrics::put_stage_timer();
|
||||
if let Err(err) = os::fsync_dir_with_owner(dir, Some(mutation_lease.clone())).await {
|
||||
if let Err(err) = os::fsync_dir_with_owner(dir, Some(commit.mutation_lease.clone())).await {
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_ANCESTOR_DIR_FSYNC,
|
||||
fsync_started,
|
||||
);
|
||||
rollback_committed_rename_std(&dst_file_path, committed_new_data_path, rollback_data_dir)
|
||||
.map_err(to_file_error)?;
|
||||
commit.rollback_non_inline_std(rollback_data_dir).map_err(to_file_error)?;
|
||||
// Same post-commit rollback window as above — drop cached
|
||||
// dst part fds so readers re-open the restored inode
|
||||
// (rustfs/backlog#1177).
|
||||
@@ -891,7 +864,7 @@ impl LocalDisk {
|
||||
// complete. Do not retain the Windows object identity guard while
|
||||
// cleaning staging paths or invalidating cached descriptors.
|
||||
#[cfg(windows)]
|
||||
drop(rename_commit_guard);
|
||||
drop(commit.directory_guard);
|
||||
|
||||
if let Some(src_file_path_parent) = src_file_path.parent() {
|
||||
if src_volume != super::super::RUSTFS_META_MULTIPART_BUCKET {
|
||||
@@ -939,7 +912,7 @@ impl LocalDisk {
|
||||
let dst_path_for_failpoint = dst_path.to_string();
|
||||
#[cfg(windows)]
|
||||
let source_parent = src_file_parent.to_path_buf();
|
||||
let rename_commit_guard_for_preparation = rename_commit_guard.clone();
|
||||
let rename_commit_guard_for_preparation = commit.directory_guard.clone();
|
||||
let sync = durability.syncs_commit_metadata();
|
||||
#[cfg(test)]
|
||||
run_inline_before_file_sync_admission(dst_path);
|
||||
@@ -1022,7 +995,7 @@ impl LocalDisk {
|
||||
if sync {
|
||||
std::fs::File::open(&backup_path)?.sync_data()?;
|
||||
}
|
||||
staged_rollback_path = Some(backup_path);
|
||||
staged_rollback_path = Some(StagedRollbackBackup(backup_path));
|
||||
}
|
||||
|
||||
Ok::<_, std::io::Error>((
|
||||
@@ -1036,9 +1009,10 @@ impl LocalDisk {
|
||||
))
|
||||
};
|
||||
let inline_preparation = if let Some(admission) = file_sync_admission.as_ref() {
|
||||
os::run_blocking_namespace_file_sync_operation(mutation_lease.clone(), admission, prepare_inline_metadata).await
|
||||
os::run_blocking_namespace_file_sync_operation(commit.mutation_lease.clone(), admission, prepare_inline_metadata)
|
||||
.await
|
||||
} else {
|
||||
os::run_blocking_namespace_operation(mutation_lease.clone(), prepare_inline_metadata).await
|
||||
os::run_blocking_namespace_operation(commit.mutation_lease.clone(), prepare_inline_metadata).await
|
||||
}
|
||||
.map_err(to_file_error)
|
||||
.map_err(DiskError::from);
|
||||
@@ -1061,28 +1035,30 @@ impl LocalDisk {
|
||||
}
|
||||
};
|
||||
|
||||
let rename_commit_guard = remove_dst_base_before_commit(
|
||||
commit.directory_guard = remove_dst_base_before_commit(
|
||||
dst_path,
|
||||
rename_commit_guard,
|
||||
commit.directory_guard,
|
||||
src_file_parent,
|
||||
dst_file_parent,
|
||||
&dst_volume_dir,
|
||||
&self.publication_root,
|
||||
mutation_lease.clone(),
|
||||
commit.publication_root,
|
||||
commit.mutation_lease.clone(),
|
||||
)
|
||||
.await?;
|
||||
|
||||
if should_remove_staged_meta_before_commit(dst_path) {
|
||||
drop(prepared_metadata_source);
|
||||
let remove_result = std::fs::remove_file(&src_file_path);
|
||||
if let Some(backup_path) = local_rollback_path.as_deref() {
|
||||
if let Some(backup_path) = local_rollback_path.as_ref().map(StagedRollbackBackup::path) {
|
||||
let _ = remove_file_if_exists(backup_path);
|
||||
}
|
||||
remove_result.map_err(to_file_error)?;
|
||||
return Err(DiskError::FileNotFound);
|
||||
}
|
||||
|
||||
if let (Some(rollback_data_dir), Some(staged_backup)) = (rollback_data_dir, local_rollback_path.as_deref()) {
|
||||
if let (Some(rollback_data_dir), Some(staged_backup)) =
|
||||
(rollback_data_dir, local_rollback_path.as_ref().map(StagedRollbackBackup::path))
|
||||
{
|
||||
let Some(dst_parent) = dst_file_path.parent() else {
|
||||
return Err(DiskError::other("missing object metadata parent"));
|
||||
};
|
||||
@@ -1096,8 +1072,8 @@ impl LocalDisk {
|
||||
staged_backup,
|
||||
&backup_path,
|
||||
&dst_volume_dir,
|
||||
&self.publication_root,
|
||||
Some(mutation_lease.clone()),
|
||||
commit.publication_root,
|
||||
Some(commit.mutation_lease.clone()),
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -1119,7 +1095,8 @@ impl LocalDisk {
|
||||
{
|
||||
let fsync_started = rustfs_io_metrics::put_stage_timer();
|
||||
if let Err(err) =
|
||||
os::fsync_dir_with_namespace_file_sync_limit(backup_parent, mutation_lease.clone(), admission).await
|
||||
os::fsync_dir_with_namespace_file_sync_limit(backup_parent, commit.mutation_lease.clone(), admission)
|
||||
.await
|
||||
{
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_BACKUP_DIR_FSYNC,
|
||||
@@ -1143,14 +1120,14 @@ impl LocalDisk {
|
||||
&src_file_path,
|
||||
&dst_file_path,
|
||||
&dst_volume_dir,
|
||||
&self.publication_root,
|
||||
&rename_commit_guard,
|
||||
mutation_lease.clone(),
|
||||
commit.publication_root,
|
||||
&commit.directory_guard,
|
||||
commit.mutation_lease.clone(),
|
||||
)
|
||||
.await
|
||||
};
|
||||
if let Err(err) = commit_result {
|
||||
if let Some(backup_path) = local_rollback_path.as_deref() {
|
||||
if let Some(backup_path) = local_rollback_path.as_ref().map(StagedRollbackBackup::path) {
|
||||
let _ = remove_file_if_exists(backup_path);
|
||||
}
|
||||
for part_path in &invalidate_part_paths {
|
||||
@@ -1161,7 +1138,7 @@ impl LocalDisk {
|
||||
|
||||
let post_commit = async {
|
||||
if should_fail_after_metadata_commit(dst_path) {
|
||||
rollback_inline_metadata_commit_std(&dst_file_path, rollback_data_dir, local_rollback_path.as_deref())?;
|
||||
commit.rollback_inline_std(rollback_data_dir, local_rollback_path.as_ref())?;
|
||||
return Err(std::io::Error::other("test fail after metadata commit"));
|
||||
}
|
||||
|
||||
@@ -1170,15 +1147,18 @@ impl LocalDisk {
|
||||
&& let Some(dst_parent) = dst_file_path.parent()
|
||||
{
|
||||
let fsync_started = rustfs_io_metrics::put_stage_timer();
|
||||
if let Err(err) =
|
||||
os::fsync_dst_dir_group_commit_or_namespace_file_sync_limit(dst_parent, mutation_lease.clone(), admission)
|
||||
.await
|
||||
if let Err(err) = os::fsync_dst_dir_group_commit_or_namespace_file_sync_limit(
|
||||
dst_parent,
|
||||
commit.mutation_lease.clone(),
|
||||
admission,
|
||||
)
|
||||
.await
|
||||
{
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_DST_DIR_FSYNC,
|
||||
fsync_started,
|
||||
);
|
||||
rollback_inline_metadata_commit_std(&dst_file_path, rollback_data_dir, local_rollback_path.as_deref())?;
|
||||
commit.rollback_inline_std(rollback_data_dir, local_rollback_path.as_ref())?;
|
||||
return Err(err);
|
||||
}
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
@@ -1204,17 +1184,14 @@ impl LocalDisk {
|
||||
}
|
||||
let fsync_started = rustfs_io_metrics::put_stage_timer();
|
||||
if let Err(err) =
|
||||
os::fsync_dir_with_namespace_file_sync_limit(ancestor_dir, mutation_lease.clone(), admission).await
|
||||
os::fsync_dir_with_namespace_file_sync_limit(ancestor_dir, commit.mutation_lease.clone(), admission)
|
||||
.await
|
||||
{
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_ANCESTOR_DIR_FSYNC,
|
||||
fsync_started,
|
||||
);
|
||||
rollback_inline_metadata_commit_std(
|
||||
&dst_file_path,
|
||||
rollback_data_dir,
|
||||
local_rollback_path.as_deref(),
|
||||
)?;
|
||||
commit.rollback_inline_std(rollback_data_dir, local_rollback_path.as_ref())?;
|
||||
return Err(err);
|
||||
}
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from(
|
||||
@@ -1251,9 +1228,9 @@ impl LocalDisk {
|
||||
// The commit no longer has a rollback path. Release the Windows
|
||||
// object identity guard before best-effort staging cleanup.
|
||||
#[cfg(windows)]
|
||||
drop(rename_commit_guard);
|
||||
drop(commit.directory_guard);
|
||||
|
||||
if let Some(backup_path) = local_rollback_path.as_deref() {
|
||||
if let Some(backup_path) = local_rollback_path.as_ref().map(StagedRollbackBackup::path) {
|
||||
let _ = remove_file_if_exists(backup_path);
|
||||
}
|
||||
|
||||
|
||||
@@ -82,7 +82,7 @@ A single filesystem rename does not atomically commit a sidecar plus `xl.meta`.
|
||||
|
||||
A laggard need not contain the predecessor. Recovery fetches the chosen decision and its verified metadata, reconstructs or validates its shards under [erasure-coding.md](erasure-coding.md), and installs that state. An empty replacement uses a fresh disk incarnation and cannot reuse old preparation receipts. A disk whose durable state claims a later decision than the supplied one rejects the operation; a conflicting same-revision digest is quarantined. Never erase a divergent disk merely because it is in the minority.
|
||||
|
||||
`rollback_committed_rename_std` and `rollback_inline_metadata_commit_std` in `crates/ecstore/src/disk/local/commit.rs`, and `restore_metadata_backup` in `crates/ecstore/src/disk/local.rs`, must become decision-aware before strict mode includes them. The permitted rollback is limited to a transaction's unchosen private preparation, or restoration proven by recovery to be necessary before any newer local effect. A chosen operation is repaired forward. Neither a client timeout nor a rename-tail error permits reverting an acknowledged decision.
|
||||
`RenameCommitContext::rollback_non_inline_std` and `RenameCommitContext::rollback_inline_std` in `crates/ecstore/src/disk/local/commit.rs`, and `restore_metadata_backup` in `crates/ecstore/src/disk/local.rs`, must become decision-aware before strict mode includes them. The permitted rollback is limited to a transaction's unchosen private preparation, or restoration proven by recovery to be necessary before any newer local effect. A chosen operation is repaired forward. Neither a client timeout nor a rename-tail error permits reverting an acknowledged decision.
|
||||
|
||||
`RenameConvergence` remains a post-publication repair signal. `PartialCommit`, `SignatureDivergent`, and `Unknown` do not decide which transaction won. Keep their diagnostics and quorum accounting; resolve authority first. Early ACK may still precede minority-tail completion after the two quorum conditions hold. Tests that inspect all disks must synchronize the tail or assert the permitted minority residue separately.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user