From 42433584abdbdf4b4b6127f25d9644ef99498117 Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 11 Aug 2026 21:05:27 +0800 Subject: [PATCH] perf(ecstore): bound strict inline commit syncs (#5931) * perf(ecstore): bound strict inline commit syncs * test(ecstore): fix admission assertion spelling * fix(ecstore): address strict inline sync review * test(ecstore): use io path for fsync hook --------- Co-authored-by: houseme --- crates/ecstore/src/disk/local.rs | 292 +++++++++++++++++++++++++++++-- crates/ecstore/src/disk/os.rs | 204 +++++++++++++++++++-- 2 files changed, 468 insertions(+), 28 deletions(-) diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 92cdb9921..47c551d79 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -2027,6 +2027,9 @@ type InlinePreparationHook = Box; static INLINE_PREPARATION_BEFORE_BACKUP: std::sync::LazyLock>> = std::sync::LazyLock::new(|| std::sync::Mutex::new(HashMap::new())); #[cfg(test)] +static INLINE_BEFORE_FILE_SYNC_ADMISSION: std::sync::LazyLock>> = + std::sync::LazyLock::new(|| std::sync::Mutex::new(HashMap::new())); +#[cfg(test)] static RENAME_DATA_AFTER_FIRST_PUBLICATION: std::sync::LazyLock>> = std::sync::LazyLock::new(|| std::sync::Mutex::new(HashMap::new())); #[cfg(test)] @@ -2091,6 +2094,14 @@ fn set_inline_preparation_before_backup(dst_path: &str, hook: impl FnOnce() + Se .insert(dst_path.to_string(), Box::new(hook)); } +#[cfg(test)] +fn set_inline_before_file_sync_admission(dst_path: &str, hook: impl FnOnce() + Send + 'static) { + INLINE_BEFORE_FILE_SYNC_ADMISSION + .lock() + .expect("test admission hook lock should not be poisoned") + .insert(dst_path.to_string(), Box::new(hook)); +} + #[cfg(test)] fn set_rename_data_after_first_publication(dst_path: &str, hook: impl FnOnce() + Send + 'static) { RENAME_DATA_AFTER_FIRST_PUBLICATION @@ -2236,6 +2247,17 @@ fn run_inline_preparation_before_backup(dst_path: &str) { } } +#[cfg(test)] +fn run_inline_before_file_sync_admission(dst_path: &str) { + let hook = INLINE_BEFORE_FILE_SYNC_ADMISSION + .lock() + .expect("test admission hook lock should not be poisoned") + .remove(dst_path); + if let Some(hook) = hook { + hook(); + } +} + #[cfg(test)] fn run_rename_data_after_first_publication(dst_path: &str) { let hook = RENAME_DATA_AFTER_FIRST_PUBLICATION @@ -9021,7 +9043,20 @@ impl DiskAPI for LocalDisk { #[cfg(windows)] let source_parent = src_file_parent.to_path_buf(); let rename_commit_guard_for_preparation = rename_commit_guard.clone(); - let inline_preparation = os::run_blocking_namespace_operation(mutation_lease.clone(), move || { + let sync = durability.syncs_commit_metadata(); + #[cfg(test)] + run_inline_before_file_sync_admission(dst_path); + let mut file_sync_admission = if sync { + Some( + os::acquire_file_sync_admission(self.file_sync_permits.clone()) + .await + .map_err(to_file_error) + .map_err(DiskError::from)?, + ) + } else { + None + }; + let prepare_inline_metadata = move || { let mut prepared_metadata_source = os::create_prepared_rename_source_with_commit_guard(&src, &dst, &rename_commit_guard_for_preparation)?; #[cfg(windows)] @@ -9058,7 +9093,6 @@ impl DiskAPI for LocalDisk { None } }); - let sync = durability.syncs_commit_metadata(); let mut staged_rollback_path = None; if let Some(d) = old_data_dir.as_ref() { let _ = xlmeta.data.remove_two(version_id, *d); @@ -9103,8 +9137,12 @@ impl DiskAPI for LocalDisk { has_dst_buf.is_none(), prepared_metadata_source, )) - }) - .await + }; + 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 + } else { + os::run_blocking_namespace_operation(mutation_lease.clone(), prepare_inline_metadata).await + } .map_err(to_file_error) .map_err(DiskError::from); @@ -9154,14 +9192,26 @@ impl DiskAPI for LocalDisk { let backup_path = dst_parent .join(rollback_data_dir.to_string()) .join(STORAGE_FORMAT_FILE_BACKUP); + // rename_all acquires the backup path's namespace lease. Do not + // hold a disk admission while acquiring another namespace lock. + drop(file_sync_admission.take()); if let Err(err) = rename_all(staged_backup, &backup_path, &dst_volume_dir, &self.publication_root).await { let _ = remove_file_if_exists(staged_backup); return Err(err); } run_rename_data_after_first_publication(dst_path); - if durability.syncs_commit_metadata() + if sync { + file_sync_admission = Some( + os::acquire_file_sync_admission(self.file_sync_permits.clone()) + .await + .map_err(to_file_error) + .map_err(DiskError::from)?, + ); + } + if let Some(admission) = file_sync_admission.as_ref() && let Some(backup_parent) = backup_path.parent() - && let Err(err) = os::fsync_dir(backup_parent).await + && let Err(err) = + os::fsync_dir_with_namespace_file_sync_limit(backup_parent, mutation_lease.clone(), admission).await { return Err(DiskError::from(to_file_error(err))); } @@ -9199,9 +9249,10 @@ impl DiskAPI for LocalDisk { } // Persist the commit rename's directory entry across power loss. - if durability.syncs_commit_metadata() + if let Some(admission) = file_sync_admission.as_ref() && let Some(dst_parent) = dst_file_path.parent() - && let Err(err) = os::fsync_dir(dst_parent).await + && let Err(err) = + os::fsync_dir_with_namespace_file_sync_limit(dst_parent, mutation_lease.clone(), admission).await { rollback_inline_metadata_commit_std(&dst_file_path, rollback_data_dir, local_rollback_path.as_deref())?; return Err(err); @@ -9214,13 +9265,17 @@ impl DiskAPI for LocalDisk { // not its own entry, so for a new inline object fsync the ancestor // chain up to and including the bucket. Overwrites already have a // durable object dir; the starts_with guard bounds the walk. - if durability.syncs_commit_metadata() && destination_was_absent { + if let Some(admission) = file_sync_admission.as_ref() + && destination_was_absent + { let mut ancestor = dst_file_path.parent().and_then(|object_dir| object_dir.parent()); while let Some(ancestor_dir) = ancestor { if !ancestor_dir.starts_with(&dst_volume_dir) { break; } - if let Err(err) = os::fsync_dir(ancestor_dir).await { + if let Err(err) = + os::fsync_dir_with_namespace_file_sync_limit(ancestor_dir, mutation_lease.clone(), admission).await + { rollback_inline_metadata_commit_std( &dst_file_path, rollback_data_dir, @@ -9239,6 +9294,10 @@ impl DiskAPI for LocalDisk { } .await; + // The disk admission protects the durability chain, not staging + // cleanup or cache invalidation after that chain has completed. + drop(file_sync_admission.take()); + // A post-commit rollback (for example, a commit-metadata fsync // failure under strict durability) restores the old metadata; drop any // descriptors cached during the committed window before propagating the @@ -12405,7 +12464,7 @@ mod test { } #[tokio::test] - async fn test_rename_data_new_inline_object_fsyncs_new_ancestor_dirs() { + async fn windows_and_unix_rename_data_new_inline_object_fsyncs_new_ancestor_dirs() { // The inline commit path (fi.data present) has the same mkdir gap as the // non-inline path: a first PUT under a new prefix must fsync the newly // created prefix and bucket dirs. @@ -12434,10 +12493,122 @@ mod test { os::fsync_dir_recorder::was_fsynced(&prefix_dir), "the newly created prefix dir must be fsynced on an inline first PUT" ); + assert_eq!( + os::fsync_dir_recorder::was_limited(&prefix_dir), + cfg!(unix), + "only Unix inline prefix fsyncs should use the disk file-sync limit" + ); assert!( os::fsync_dir_recorder::was_fsynced(&bucket_dir), "the bucket dir must be fsynced on an inline first PUT" ); + assert_eq!( + os::fsync_dir_recorder::was_limited(&bucket_dir), + cfg!(unix), + "only Unix inline bucket fsyncs should use the disk file-sync limit" + ); + } + + #[cfg(unix)] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[allow(clippy::await_holding_lock)] + async fn strict_inline_rename_retains_admission_until_commit_fsync() { + use std::sync::atomic::{AtomicU8, Ordering}; + use std::sync::mpsc; + use tempfile::tempdir; + use tokio::sync::oneshot; + + const FIRST_BARRIER: u8 = 1; + const SECOND_PREPARATION: u8 = 2; + + let _mode = durability_mode_override::set(DurabilityMode::Strict); + 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 mut disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + disk.file_sync_permits = Arc::new(Semaphore::new(1)); + let disk = Arc::new(disk); + let bucket = "inline-admission-order"; + let first_object = "first-object"; + let second_object = "second-object"; + ensure_test_volume(&disk, bucket).await; + ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + + let (first_prepared_tx, first_prepared_rx) = mpsc::channel(); + let (release_first_tx, release_first_rx) = mpsc::channel(); + set_inline_preparation_before_backup(first_object, move || { + first_prepared_tx.send(()).expect("signal first preparation"); + release_first_rx.recv().expect("wait for queued rename"); + }); + let first_disk = disk.clone(); + let first = tokio::spawn(async move { + first_disk + .rename_data( + RUSTFS_META_TMP_BUCKET, + "first-stage", + test_file_info(first_object, Uuid::new_v4(), None, Some(Bytes::from_static(b"first"))), + bucket, + first_object, + ) + .await + }); + tokio::task::spawn_blocking(move || first_prepared_rx.recv_timeout(Duration::from_secs(30))) + .await + .expect("first preparation waiter should run") + .expect("first rename should hold the only disk admission"); + + let first_event = Arc::new(AtomicU8::new(0)); + let first_barrier_event = first_event.clone(); + let first_object_dir = disk + .get_object_path_for_io(bucket, first_object) + .expect("first object path should resolve"); + os::fsync_dir_recorder::set_before_limited(&first_object_dir, move || { + let _ = first_barrier_event.compare_exchange(0, FIRST_BARRIER, Ordering::SeqCst, Ordering::SeqCst); + }); + + let second_preparation_event = first_event.clone(); + set_inline_preparation_before_backup(second_object, move || { + second_preparation_event.fetch_or(SECOND_PREPARATION, Ordering::SeqCst); + }); + let (second_admission_tx, second_admission_rx) = oneshot::channel(); + set_inline_before_file_sync_admission(second_object, move || { + second_admission_tx.send(()).expect("signal second admission attempt"); + }); + let second_disk = disk.clone(); + let mut second = Box::pin(async move { + second_disk + .rename_data( + RUSTFS_META_TMP_BUCKET, + "second-stage", + test_file_info(second_object, Uuid::new_v4(), None, Some(Bytes::from_static(b"second"))), + bucket, + second_object, + ) + .await + }); + let mut second_admission_rx = Box::pin(second_admission_rx); + tokio::time::timeout(Duration::from_secs(30), async { + tokio::select! { + _ = &mut second => panic!("second rename must wait for disk admission"), + signal = &mut second_admission_rx => signal.expect("second admission hook should run"), + } + }) + .await + .expect("second rename should reach the admission queue"); + + release_first_tx.send(()).expect("release first preparation"); + let (first_result, second_result) = tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(first, second) }) + .await + .expect("both inline renames should complete"); + first_result + .expect("first rename task should join") + .expect("first inline rename should commit"); + second_result.expect("second inline rename should commit"); + + assert_eq!( + first_event.load(Ordering::SeqCst), + FIRST_BARRIER | SECOND_PREPARATION, + "the queued rename must not overtake the admitted rename before its commit fsync" + ); } #[cfg(windows)] @@ -13684,14 +13855,16 @@ mod test { ); } - #[tokio::test] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[allow(clippy::await_holding_lock)] async fn test_rename_data_writes_old_metadata_backup_for_inline_overwrite() { + use std::sync::mpsc; use tempfile::tempdir; let _mode = durability_mode_override::set(DurabilityMode::Strict); 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 disk = Arc::new(LocalDisk::new(&endpoint, false).await.expect("local disk should be created")); let bucket = "bucket"; let object = "inline-object"; @@ -13717,10 +13890,32 @@ mod test { .await .expect("tmp object dir should be created"); + let (published_tx, published_rx) = mpsc::channel(); + let (release_tx, release_rx) = mpsc::channel(); + set_rename_data_after_first_publication(object, move || { + published_tx.send(()).expect("signal backup publication"); + release_rx.recv().expect("wait for lock-order assertion"); + }); let new_fi = test_file_info(object, version_id, None, Some(Bytes::from_static(b"inline-new"))); - let resp = disk - .rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) + let rename_disk = disk.clone(); + let rename = tokio::spawn(async move { + rename_disk + .rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) + .await + }); + tokio::task::spawn_blocking(move || published_rx.recv_timeout(Duration::from_secs(10))) .await + .expect("publication waiter should run") + .expect("rollback backup must be published"); + assert_eq!( + disk.file_sync_permits.available_permits(), + os::MAX_PARALLEL_FILE_SYNCS, + "backup publication must not acquire namespace while holding disk admission" + ); + release_tx.send(()).expect("release backup publication"); + let resp = rename + .await + .expect("inline rename task should join") .expect("inline rename_data should commit"); assert_eq!(resp.old_data_dir, Some(old_data_dir)); @@ -13731,10 +13926,20 @@ mod test { os::fsync_dir_recorder::was_fsynced(backup_path.parent().expect("backup must have a parent")), "strict inline overwrite must persist the rollback backup directory entry" ); + assert_eq!( + os::fsync_dir_recorder::was_limited(backup_path.parent().expect("backup must have a parent")), + cfg!(unix), + "only Unix rollback backup fsyncs should use the disk file-sync limit" + ); assert!( os::fsync_dir_recorder::was_fsynced(&dst_object_dir), "strict inline overwrite must persist the committed xl.meta directory entry" ); + assert_eq!( + os::fsync_dir_recorder::was_limited(&dst_object_dir), + cfg!(unix), + "only Unix inline commit fsyncs should use the disk file-sync limit" + ); // The rollback backup must contain the previous metadata bytes verbatim so // that undo_write can restore the prior committed object; guards the inline // backup write against truncation/corruption regressions. @@ -14331,10 +14536,12 @@ mod test { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[allow(clippy::await_holding_lock)] async fn windows_and_unix_cancelled_inline_preparation_serializes_newer_commit() { use std::sync::mpsc; use tempfile::tempdir; + let _mode = durability_mode_override::set(DurabilityMode::Strict); 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 = Arc::new(LocalDisk::new(&endpoint, false).await.expect("local disk should be created")); @@ -14373,6 +14580,11 @@ mod test { .await .expect("preparation waiter should run") .expect("preparation must reach the backup hook"); + assert_eq!( + disk.file_sync_permits.available_permits(), + os::MAX_PARALLEL_FILE_SYNCS - 1, + "strict inline preparation must hold one disk file-sync permit" + ); cancelled.abort(); assert!(cancelled.await.expect_err("operation should be cancelled").is_cancelled()); @@ -14446,6 +14658,56 @@ mod test { ); } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[allow(clippy::await_holding_lock)] + async fn relaxed_inline_preparation_does_not_use_file_sync_limit() { + use std::sync::mpsc; + use tempfile::tempdir; + + let _mode = durability_mode_override::set(DurabilityMode::Relaxed); + 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 = Arc::new(LocalDisk::new(&endpoint, false).await.expect("local disk should be created")); + let bucket = "relaxed-inline-preparation"; + let object = "inline-object"; + ensure_test_volume(&disk, bucket).await; + ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + + let (entered_tx, entered_rx) = mpsc::channel(); + let (release_tx, release_rx) = mpsc::channel(); + set_inline_preparation_before_backup(object, move || { + entered_tx.send(()).expect("signal blocked preparation"); + release_rx.recv().expect("wait for permit assertion"); + }); + let rename_disk = Arc::clone(&disk); + let rename = tokio::spawn(async move { + rename_disk + .rename_data( + RUSTFS_META_TMP_BUCKET, + "relaxed-inline-stage", + test_file_info(object, Uuid::new_v4(), None, Some(Bytes::from_static(b"payload"))), + bucket, + object, + ) + .await + }); + tokio::task::spawn_blocking(move || entered_rx.recv_timeout(Duration::from_secs(10))) + .await + .expect("preparation waiter should run") + .expect("preparation must reach the hook"); + assert_eq!( + disk.file_sync_permits.available_permits(), + os::MAX_PARALLEL_FILE_SYNCS, + "relaxed inline preparation must not consume strict sync capacity" + ); + + release_tx.send(()).expect("release inline preparation"); + rename + .await + .expect("rename task should join") + .expect("relaxed inline rename should commit"); + } + #[tokio::test] async fn rename_purge_pending_payload_stays_object_and_cleans_local_backup() { use tempfile::tempdir; diff --git a/crates/ecstore/src/disk/os.rs b/crates/ecstore/src/disk/os.rs index f8f87975a..bd3ec77a5 100644 --- a/crates/ecstore/src/disk/os.rs +++ b/crates/ecstore/src/disk/os.rs @@ -78,28 +78,59 @@ pub fn check_path_length(path_name: &str) -> Result<()> { /// their own unique tempdir to stay robust against parallel test execution. #[cfg(test)] pub(crate) mod fsync_dir_recorder { + use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Mutex; - static RECORDED: Mutex> = Mutex::new(Vec::new()); + type Hook = Box; - pub(crate) fn record(dir: &Path) { - let mut recorded = RECORDED.lock().expect("fsync dir recorder poisoned"); - recorded.push(dir.to_path_buf()); - if let Ok(canonical) = dir.canonicalize() - && canonical != dir + static RECORDED: Mutex> = Mutex::new(Vec::new()); + static LIMITED: Mutex> = Mutex::new(Vec::new()); + static BEFORE_LIMITED: std::sync::LazyLock>> = + std::sync::LazyLock::new(|| Mutex::new(HashMap::new())); + + fn record_path(paths: &Mutex>, path: &Path, description: &str) { + let mut paths = paths.lock().expect(description); + paths.push(path.to_path_buf()); + if let Ok(canonical) = path.canonicalize() + && canonical != path { - recorded.push(canonical); + paths.push(canonical); } } - pub(crate) fn was_fsynced(dir: &Path) -> bool { - let canonical = dir.canonicalize().ok(); - RECORDED - .lock() - .expect("fsync dir recorder poisoned") + fn contains_path(paths: &[PathBuf], path: &Path) -> bool { + let canonical = path.canonicalize().ok(); + paths .iter() - .any(|p| p == dir || canonical.as_ref().is_some_and(|canonical| p == canonical)) + .any(|recorded| recorded == path || canonical.as_ref().is_some_and(|canonical| recorded == canonical)) + } + + pub(crate) fn record(dir: &Path) { + record_path(&RECORDED, dir, "fsync dir recorder"); + } + + pub(crate) fn was_fsynced(dir: &Path) -> bool { + contains_path(&RECORDED.lock().expect("fsync dir recorder poisoned"), dir) + } + + pub(crate) fn record_limited(dir: &Path) { + record_path(&LIMITED, dir, "limited fsync dir recorder"); + let hook = BEFORE_LIMITED.lock().expect("limited fsync hook poisoned").remove(dir); + if let Some(hook) = hook { + hook(); + } + } + + pub(crate) fn was_limited(dir: &Path) -> bool { + contains_path(&LIMITED.lock().expect("limited fsync dir recorder poisoned"), dir) + } + + pub(crate) fn set_before_limited(dir: &Path, hook: impl FnOnce() + Send + 'static) { + BEFORE_LIMITED + .lock() + .expect("limited fsync hook poisoned") + .insert(dir.to_path_buf(), Box::new(hook)); } } @@ -1120,6 +1151,79 @@ pub(crate) async fn run_blocking_namespace_operation( .map_err(|err| io::Error::other(format!("blocking namespace operation failed: {err}")))? } +/// Admit one strict inline commit under the disk sync limit. The caller already +/// owns the namespace lease, establishing namespace -> disk ordering. Holding +/// admission across adjacent durability barriers prevents one transaction from +/// repeatedly joining the disk semaphore tail. +pub(crate) struct FileSyncAdmission { + disk_permit: Arc, +} + +pub(crate) async fn acquire_file_sync_admission(disk_permits: Arc) -> io::Result { + let disk_permit = disk_permits + .acquire_owned() + .await + .map_err(|_| io::Error::other("disk file sync concurrency limiter closed"))?; + Ok(FileSyncAdmission { + disk_permit: Arc::new(disk_permit), + }) +} + +/// Keep the disk admission and namespace lease with the blocking syscall if +/// the async waiter is cancelled. The process-wide admission remains with the +/// waiter so cancellation cannot starve healthy disks. +pub(crate) async fn run_blocking_namespace_file_sync_operation( + lease: Arc, + admission: &FileSyncAdmission, + operation: impl FnOnce() -> io::Result + Send + 'static, +) -> io::Result { + run_blocking_namespace_file_sync_operation_with_global(lease, admission, &FILE_SYNC_PERMITS, operation).await +} + +async fn run_blocking_namespace_file_sync_operation_with_global( + lease: Arc, + admission: &FileSyncAdmission, + global_permits: &Semaphore, + operation: impl FnOnce() -> io::Result + Send + 'static, +) -> io::Result { + let global_permit = global_permits + .acquire() + .await + .map_err(|_| io::Error::other("global file sync concurrency limiter closed"))?; + let disk_permit = admission.disk_permit.clone(); + let result = tokio::task::spawn_blocking(move || { + let _lease = lease; + let _disk_permit = disk_permit; + operation() + }) + .await; + drop(global_permit); + result.map_err(|err| io::Error::other(format!("blocking namespace file sync operation failed: {err}")))? +} + +pub(crate) async fn fsync_dir_with_namespace_file_sync_limit( + dir: impl AsRef, + lease: Arc, + admission: &FileSyncAdmission, +) -> io::Result<()> { + #[cfg(unix)] + { + let dir = dir.as_ref().to_path_buf(); + run_blocking_namespace_file_sync_operation(lease, admission, move || { + #[cfg(test)] + fsync_dir_recorder::record_limited(&dir); + fsync_dir_std(dir) + }) + .await + } + + #[cfg(not(unix))] + { + let _ = (lease, admission); + fsync_dir_std(dir) + } +} + struct RenamePreparation { parent_guard: Option, #[cfg(windows)] @@ -2784,6 +2888,7 @@ pub fn is_dir_not_empty_error(err: &io::Error) -> bool { mod tests { use super::*; use std::sync::Mutex; + use std::time::Duration; use tempfile::tempdir; use tracing_subscriber::fmt::MakeWriter; @@ -4553,6 +4658,79 @@ mod tests { fsync_dir(temp_dir.path()).await.expect("fsync dir must succeed"); } + #[tokio::test] + async fn file_sync_admission_is_reused_across_commit_barriers() { + let temp_dir = tempdir().expect("create temp dir"); + let limiter = Arc::new(Semaphore::new(1)); + let lease = acquire_namespace_mutation_lease(temp_dir.path()).await; + let admission = acquire_file_sync_admission(limiter.clone()) + .await + .expect("first commit should acquire admission"); + + run_blocking_namespace_file_sync_operation(lease.clone(), &admission, || Ok(())) + .await + .expect("first barrier should complete under the admission"); + let mut waiting = Box::pin(acquire_file_sync_admission(limiter)); + assert!( + futures::poll!(&mut waiting).is_pending(), + "another commit must remain queued between durability barriers" + ); + run_blocking_namespace_file_sync_operation(lease, &admission, || Ok(())) + .await + .expect("later barrier should reuse admission without requeuing"); + + drop(admission); + tokio::time::timeout(Duration::from_secs(30), waiting) + .await + .expect("queued commit should acquire admission after release") + .expect("queued commit should acquire admission"); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn cancelled_file_sync_waiter_keeps_disk_admission_until_blocking_work_finishes() { + use std::sync::mpsc; + + let temp_dir = tempdir().expect("create temp dir"); + let limiter = Arc::new(Semaphore::new(1)); + let global_permits = Arc::new(Semaphore::new(1)); + let lease = acquire_namespace_mutation_lease(temp_dir.path()).await; + let admission = acquire_file_sync_admission(limiter.clone()) + .await + .expect("file sync admission should be acquired"); + let (entered_tx, entered_rx) = mpsc::channel(); + let (release_tx, release_rx) = mpsc::channel(); + let waiter_global_permits = global_permits.clone(); + let waiter = tokio::spawn(async move { + run_blocking_namespace_file_sync_operation_with_global(lease, &admission, waiter_global_permits.as_ref(), move || { + entered_tx.send(()).expect("signal blocking work"); + release_rx.recv().expect("wait for blocking work release"); + Ok(()) + }) + .await + }); + + tokio::task::spawn_blocking(move || entered_rx.recv_timeout(Duration::from_secs(30))) + .await + .expect("blocking work waiter should run") + .expect("blocking work should start"); + waiter.abort(); + assert!(waiter.await.expect_err("waiter should be cancelled").is_cancelled()); + let returned_global_permit = global_permits + .try_acquire() + .expect("cancelled waiter must return global capacity for healthy disks"); + assert!( + limiter.clone().try_acquire_owned().is_err(), + "cancelled waiter must not return disk capacity while blocking work is active" + ); + + release_tx.send(()).expect("release blocking work"); + let _returned_permit = tokio::time::timeout(Duration::from_secs(30), limiter.acquire_owned()) + .await + .expect("disk capacity should return after blocking work finishes") + .expect("disk limiter should remain open"); + drop(returned_global_permit); + } + #[tokio::test] #[serial_test::serial(file_sync_probe)] async fn sync_dir_files_syncs_regular_files_and_dir() {