fix(ecstore): roll back failed CAS directory fsync

Restore the previous control-file bytes, or remove a newly created file, when the Unix compare-and-update path reaches the rename but then fails to fsync the parent directory. This keeps failed metadata CAS publications from advancing recovery anchors.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-07 11:12:12 +08:00
parent 672087ec0d
commit 022f1ecce5
2 changed files with 161 additions and 4 deletions
+141 -4
View File
@@ -8708,6 +8708,34 @@ impl DiskAPI for LocalDisk {
}); });
} }
let rollback_after_fsync_failure = |parent: &Path, current: Option<&[u8]>| -> std::io::Result<()> {
match current {
Some(previous) => {
let rollback_temporary =
parent.join(format!(".{}.{}.rollback.tmp", path.replace('/', "_"), Uuid::new_v4()));
let rollback_result = (|| -> std::io::Result<()> {
let mut staged = std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(&rollback_temporary)?;
staged.write_all(previous)?;
staged.sync_all()?;
std::fs::rename(&rollback_temporary, &file_path)
})();
if let Err(err) = rollback_result {
let _ = std::fs::remove_file(&rollback_temporary);
return Err(err);
}
}
None => match std::fs::remove_file(&file_path) {
Ok(()) => {}
Err(err) if err.kind() == ErrorKind::NotFound => {}
Err(err) => return Err(err),
},
}
os::fsync_dir_std(parent)
};
match replacement { match replacement {
Some(replacement) => { Some(replacement) => {
let parent = file_path let parent = file_path
@@ -8727,14 +8755,19 @@ impl DiskAPI for LocalDisk {
let _ = std::fs::remove_file(&temporary); let _ = std::fs::remove_file(&temporary);
return Err(err); return Err(err);
} }
if sync_metadata { if sync_metadata && let Err(err) = os::fsync_dir_std(parent) {
os::fsync_dir_std(parent)?; rollback_after_fsync_failure(parent, current.as_deref())?;
return Err(err);
} }
} }
None => { None => {
std::fs::remove_file(&file_path)?; std::fs::remove_file(&file_path)?;
if sync_metadata && let Some(parent) = file_path.parent() { if sync_metadata
os::fsync_dir_std(parent)?; && let Some(parent) = file_path.parent()
&& let Err(err) = os::fsync_dir_std(parent)
{
rollback_after_fsync_failure(parent, current.as_deref())?;
return Err(err);
} }
} }
} }
@@ -22166,6 +22199,110 @@ mod test {
)); ));
} }
#[cfg(unix)]
#[tokio::test]
async fn conditional_file_update_dir_fsync_failure_restores_previous_bytes() {
use tempfile::tempdir;
let dir = 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");
let previous = Bytes::from_static(b"previous-owner");
let successor = Bytes::from_static(b"successor-owner");
assert_eq!(
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, None, Some(previous.clone()))
.await
.expect("previous owner should commit"),
ConditionalFileUpdate::Updated
);
let marker_path = disk
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.expect("marker path should resolve");
let parent = marker_path.parent().expect("marker path should have a parent");
os::fsync_dir_recorder::set_failure(parent, ErrorKind::Other);
let err = disk
.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, Some(previous.clone()), Some(successor))
.await
.expect_err("directory fsync failure must fail the CAS update");
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other));
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.await
.expect("previous bytes should remain readable after rollback"),
previous
);
}
#[cfg(unix)]
#[tokio::test]
async fn conditional_file_update_dir_fsync_failure_removes_new_file_without_anchor() {
use tempfile::tempdir;
let dir = 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 parent = marker_path.parent().expect("marker path should have a parent");
os::fsync_dir_recorder::set_failure(parent, ErrorKind::Other);
let err = disk
.compare_and_update_file(
RUSTFS_META_BUCKET,
HEALING_MARKER_PATH,
None,
Some(Bytes::from_static(b"successor-owner")),
)
.await
.expect_err("directory fsync failure must fail the CAS create");
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other));
assert!(
matches!(disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH).await, Err(DiskError::FileNotFound)),
"uncommitted successor bytes must be removed when no previous anchor exists"
);
}
#[cfg(unix)]
#[tokio::test]
async fn conditional_file_delete_dir_fsync_failure_restores_previous_bytes() {
use tempfile::tempdir;
let dir = 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");
let previous = Bytes::from_static(b"previous-owner");
assert_eq!(
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, None, Some(previous.clone()))
.await
.expect("previous owner should commit"),
ConditionalFileUpdate::Updated
);
let marker_path = disk
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.expect("marker path should resolve");
let parent = marker_path.parent().expect("marker path should have a parent");
os::fsync_dir_recorder::set_failure(parent, ErrorKind::Other);
let err = disk
.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, Some(previous.clone()), None)
.await
.expect_err("directory fsync failure must fail the CAS delete");
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other));
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.await
.expect("previous bytes should be restored after failed delete"),
previous
);
}
#[cfg(unix)] #[cfg(unix)]
#[tokio::test] #[tokio::test]
async fn conditional_file_update_returns_would_block_when_marker_lock_is_contended() { async fn conditional_file_update_returns_would_block_when_marker_lock_is_contended() {
+20
View File
@@ -92,6 +92,9 @@ pub(crate) mod fsync_dir_recorder {
static LIMITED: Mutex<Vec<PathBuf>> = Mutex::new(Vec::new()); static LIMITED: Mutex<Vec<PathBuf>> = Mutex::new(Vec::new());
static GROUPED: Mutex<Vec<(PathBuf, usize)>> = Mutex::new(Vec::new()); static GROUPED: Mutex<Vec<(PathBuf, usize)>> = Mutex::new(Vec::new());
#[cfg(unix)] #[cfg(unix)]
static FAILURES: std::sync::LazyLock<Mutex<HashMap<PathBuf, io::ErrorKind>>> =
std::sync::LazyLock::new(|| Mutex::new(HashMap::new()));
#[cfg(unix)]
static BEFORE_LIMITED: std::sync::LazyLock<Mutex<HashMap<PathBuf, Hook>>> = static BEFORE_LIMITED: std::sync::LazyLock<Mutex<HashMap<PathBuf, Hook>>> =
std::sync::LazyLock::new(|| Mutex::new(HashMap::new())); std::sync::LazyLock::new(|| Mutex::new(HashMap::new()));
static BEFORE_GROUP_BATCH: std::sync::LazyLock<Mutex<HashMap<PathBuf, Hook>>> = static BEFORE_GROUP_BATCH: std::sync::LazyLock<Mutex<HashMap<PathBuf, Hook>>> =
@@ -151,6 +154,19 @@ pub(crate) mod fsync_dir_recorder {
contains_path(&RECORDED.lock().expect("fsync dir recorder poisoned"), dir) contains_path(&RECORDED.lock().expect("fsync dir recorder poisoned"), dir)
} }
#[cfg(unix)]
pub(crate) fn set_failure(dir: &Path, kind: io::ErrorKind) {
FAILURES
.lock()
.expect("fsync dir failure hook poisoned")
.insert(dir.to_path_buf(), kind);
}
#[cfg(unix)]
pub(crate) fn take_failure(dir: &Path) -> Option<io::ErrorKind> {
remove_path_keyed(&FAILURES, dir, "fsync dir failure hook poisoned")
}
#[cfg(unix)] #[cfg(unix)]
pub(crate) fn record_limited(dir: &Path) { pub(crate) fn record_limited(dir: &Path) {
record_path(&LIMITED, dir, "limited fsync dir recorder"); record_path(&LIMITED, dir, "limited fsync dir recorder");
@@ -426,6 +442,10 @@ pub fn fsync_dir_std(dir: impl AsRef<Path>) -> io::Result<()> {
fsync_dir_recorder::record(dir.as_ref()); fsync_dir_recorder::record(dir.as_ref());
#[cfg(unix)] #[cfg(unix)]
{ {
#[cfg(test)]
if let Some(kind) = fsync_dir_recorder::take_failure(dir.as_ref()) {
return Err(io::Error::from(kind));
}
std::fs::File::open(dir.as_ref())?.sync_all()?; std::fs::File::open(dir.as_ref())?.sync_all()?;
} }
#[cfg(not(unix))] #[cfg(not(unix))]