mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 04:16:38 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e99deef026 | |||
| 2f0918f60b | |||
| 1a6b870eb5 |
@@ -57,6 +57,13 @@ pub const DEFAULT_MAX_IO_EVENTS_PER_TICK: usize = 1024;
|
||||
pub const DEFAULT_EVENT_INTERVAL: u32 = 61;
|
||||
pub const DEFAULT_RNG_SEED: Option<u64> = None; // None means random
|
||||
|
||||
/// Dedicated blocking thread pool for fsync/fdatasync operations.
|
||||
/// When > 1, fsync operations are isolated from the main blocking pool to
|
||||
/// prevent device-bound fsync from starving read operations (pread/stat/open).
|
||||
/// Default 0 means auto (no isolation, use main runtime).
|
||||
pub const ENV_FSYNC_BLOCKING_THREADS: &str = "RUSTFS_RUNTIME_FSYNC_BLOCKING_THREADS";
|
||||
pub const DEFAULT_FSYNC_BLOCKING_THREADS: usize = 0;
|
||||
|
||||
// Dial9 Tokio Telemetry Default values
|
||||
pub const DEFAULT_RUNTIME_DIAL9_ENABLED: bool = false; // Disabled by default
|
||||
pub const DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR: &str = "/var/log/rustfs/telemetry";
|
||||
|
||||
+309
-1211
File diff suppressed because it is too large
Load Diff
@@ -315,7 +315,7 @@ pub async fn fsync_dir(dir: impl AsRef<Path>) -> io::Result<()> {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let dir = dir.as_ref().to_path_buf();
|
||||
tokio::task::spawn_blocking(move || fsync_dir_std(dir)).await?
|
||||
fsync_spawn_blocking(move || fsync_dir_std(dir)).await?
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
@@ -683,7 +683,7 @@ async fn fsync_open_dst_dir_group(group: &DstDirFsyncGroup) -> io::Result<()> {
|
||||
#[cfg(test)]
|
||||
let dir = group.dir.clone();
|
||||
let dir_file = group.dir_file.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
fsync_spawn_blocking(move || {
|
||||
#[cfg(test)]
|
||||
{
|
||||
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
|
||||
@@ -1080,6 +1080,44 @@ const TEST_GLOBAL_FILE_SYNCS: usize = 64;
|
||||
|
||||
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
||||
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
|
||||
/// Dedicated tokio runtime for fsync/fdatasync blocking operations. When
|
||||
/// configured with >1 threads, isolates device-bound fsync from the main
|
||||
/// blocking pool so reads (pread/stat/open) are not starved. `None` means
|
||||
/// fall back to the main runtime (zero behavior change).
|
||||
static FSYNC_RUNTIME: LazyLock<Option<tokio::runtime::Runtime>> = LazyLock::new(|| {
|
||||
let threads =
|
||||
rustfs_utils::get_env_usize(rustfs_config::ENV_FSYNC_BLOCKING_THREADS, rustfs_config::DEFAULT_FSYNC_BLOCKING_THREADS);
|
||||
if threads <= 1 {
|
||||
return None;
|
||||
}
|
||||
let mut builder = tokio::runtime::Builder::new_multi_thread();
|
||||
builder
|
||||
.worker_threads(num_cpus::get().min(8))
|
||||
.max_blocking_threads(threads)
|
||||
.thread_name("rustfs-fsync")
|
||||
.thread_stack_size(512 * 1024)
|
||||
.enable_all();
|
||||
match builder.build() {
|
||||
Ok(rt) => {
|
||||
tracing::info!(threads, "fsync dedicated blocking pool enabled");
|
||||
Some(rt)
|
||||
}
|
||||
Err(err) => {
|
||||
tracing::warn!(%err, "failed to build fsync runtime, falling back to main pool");
|
||||
None
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
/// Spawn a blocking task on the fsync-dedicated runtime if configured,
|
||||
/// otherwise fall back to the main tokio blocking pool.
|
||||
fn fsync_spawn_blocking<T: Send + 'static>(f: impl FnOnce() -> T + Send + 'static) -> tokio::task::JoinHandle<T> {
|
||||
match FSYNC_RUNTIME.as_ref() {
|
||||
Some(rt) => rt.spawn_blocking(f),
|
||||
None => tokio::task::spawn_blocking(f),
|
||||
}
|
||||
}
|
||||
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
type NamespaceMutationLock = AsyncMutex<()>;
|
||||
@@ -1217,7 +1255,7 @@ where
|
||||
F: FnOnce() -> io::Result<T> + Send + 'static,
|
||||
{
|
||||
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
let result = fsync_spawn_blocking(move || {
|
||||
let _disk_permit = disk_permit;
|
||||
work()
|
||||
})
|
||||
@@ -2146,7 +2184,7 @@ async fn run_blocking_namespace_file_sync_operation_with_global<T: Send + 'stati
|
||||
wait_started,
|
||||
);
|
||||
let disk_permit = admission.disk_permit.clone();
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
let result = fsync_spawn_blocking(move || {
|
||||
let _lease = lease;
|
||||
let _disk_permit = disk_permit;
|
||||
operation()
|
||||
|
||||
@@ -160,10 +160,6 @@ pub struct InstanceContext {
|
||||
/// workers (scanner/heal/tier/lifecycle) without touching another instance.
|
||||
/// Replaces the process-global cancel-token static.
|
||||
background_cancel_token: OnceLock<CancellationToken>,
|
||||
/// Serializes decommission data-movement operations with cancellation and
|
||||
/// a subsequent restart. Readers are held across one object side effect;
|
||||
/// the transition path takes the writer after cancelling the routine.
|
||||
decommission_operation_gate: Arc<RwLock<()>>,
|
||||
/// Resolves object-encryption material at the application boundary.
|
||||
object_encryption_resolver: OnceLock<Arc<dyn ObjectEncryptionResolver>>,
|
||||
tier_delete_journal_recovery_stores: std::sync::Mutex<HashSet<Uuid>>,
|
||||
@@ -204,7 +200,6 @@ impl InstanceContext {
|
||||
local_disk_set_drives: Arc::new(RwLock::new(Vec::new())),
|
||||
bucket_metadata_sys: std::sync::Mutex::new(None),
|
||||
background_cancel_token: OnceLock::new(),
|
||||
decommission_operation_gate: Arc::new(RwLock::new(())),
|
||||
object_encryption_resolver: OnceLock::new(),
|
||||
tier_delete_journal_recovery_stores: std::sync::Mutex::new(HashSet::new()),
|
||||
transition_transaction_recovery_stores: std::sync::Mutex::new(HashSet::new()),
|
||||
@@ -223,10 +218,6 @@ impl InstanceContext {
|
||||
self.lock_manager.clone()
|
||||
}
|
||||
|
||||
pub(crate) fn decommission_operation_gate(&self) -> Arc<RwLock<()>> {
|
||||
Arc::clone(&self.decommission_operation_gate)
|
||||
}
|
||||
|
||||
/// Install the application-owned object-encryption resolver once.
|
||||
pub fn set_object_encryption_resolver(
|
||||
&self,
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -110,7 +110,10 @@ impl HealStorageAPI for MockStorage {
|
||||
Ok(Vec::new())
|
||||
}
|
||||
|
||||
async fn get_bucket_info(&self, _bucket: &str) -> Result<Option<BucketInfo>> {
|
||||
async fn get_bucket_info(&self, bucket: &str) -> Result<Option<BucketInfo>> {
|
||||
if bucket == "panic" {
|
||||
panic!("test-only panic payload must not escape the scheduler");
|
||||
}
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
@@ -1021,6 +1024,231 @@ async fn test_task_alias_is_removed_after_terminal_completion() {
|
||||
assert_eq!(manager.canonical_task_id(&duplicate_id).await, duplicate_id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn scheduler_panic_releases_active_slot_and_allows_same_target_readmission() {
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let manager = HealManager::new(storage, None);
|
||||
let request = bucket_request("panic", HealPriority::Normal, HealRequestSource::Admin);
|
||||
let task_id = request.id.clone();
|
||||
|
||||
assert_eq!(
|
||||
manager
|
||||
.submit_heal_request(request)
|
||||
.await
|
||||
.expect("panic request should be admitted"),
|
||||
HealAdmissionResult::Accepted
|
||||
);
|
||||
let duplicate = bucket_request("panic", HealPriority::Normal, HealRequestSource::Admin);
|
||||
let duplicate_id = duplicate.id.clone();
|
||||
assert_eq!(
|
||||
manager
|
||||
.submit_heal_request(duplicate)
|
||||
.await
|
||||
.expect("same target should merge while active is queued"),
|
||||
HealAdmissionResult::Merged
|
||||
);
|
||||
assert_eq!(manager.canonical_task_id(&duplicate_id).await, task_id);
|
||||
process_manager_queue_once(&manager).await;
|
||||
|
||||
let status = tokio::time::timeout(Duration::from_secs(1), async {
|
||||
loop {
|
||||
if let Ok(status) = manager.get_task_status(&task_id).await
|
||||
&& matches!(status, HealTaskStatus::Failed { .. })
|
||||
{
|
||||
break status;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("panic task should reach a terminal status");
|
||||
assert_eq!(
|
||||
status,
|
||||
HealTaskStatus::Failed {
|
||||
error: PANICKED_HEAL_TASK_ERROR.to_string()
|
||||
}
|
||||
);
|
||||
assert_eq!(manager.get_active_task_count().await, 0);
|
||||
assert_eq!(manager.get_queue_length().await, 0);
|
||||
assert!(manager.retrying_heals.lock().await.is_empty());
|
||||
assert!(manager.task_aliases.lock().await.is_empty());
|
||||
assert!(manager.completed_heals.lock().await.contains_key(&task_id));
|
||||
assert_eq!(manager.canonical_task_id(&duplicate_id).await, duplicate_id);
|
||||
|
||||
let readmitted = bucket_request("panic", HealPriority::Normal, HealRequestSource::Admin);
|
||||
assert_eq!(
|
||||
manager
|
||||
.submit_heal_request(readmitted)
|
||||
.await
|
||||
.expect("same target should be re-admitted after a panic"),
|
||||
HealAdmissionResult::Accepted
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn retry_child_panic_finishes_parent_once() {
|
||||
clear_scheduler_panic();
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let manager = HealManager::new(storage, None);
|
||||
let request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None);
|
||||
let task_id = request.id.clone();
|
||||
assert_eq!(
|
||||
manager
|
||||
.submit_heal_request(request)
|
||||
.await
|
||||
.expect("retry request should be admitted"),
|
||||
HealAdmissionResult::Accepted
|
||||
);
|
||||
arm_scheduler_panic(SchedulerPanicPoint::RetryChild, &task_id);
|
||||
process_manager_queue_once(&manager).await;
|
||||
|
||||
let status = tokio::time::timeout(Duration::from_secs(1), async {
|
||||
loop {
|
||||
if let Ok(status) = manager.get_task_status(&task_id).await
|
||||
&& matches!(status, HealTaskStatus::Failed { .. })
|
||||
{
|
||||
break status;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("retry child panic should finish the parent");
|
||||
clear_scheduler_panic();
|
||||
assert_eq!(
|
||||
status,
|
||||
HealTaskStatus::Failed {
|
||||
error: PANICKED_HEAL_TASK_ERROR.to_string()
|
||||
}
|
||||
);
|
||||
assert_eq!(manager.get_active_task_count().await, 0);
|
||||
assert_eq!(manager.get_queue_length().await, 0);
|
||||
assert!(manager.retrying_heals.lock().await.is_empty());
|
||||
assert!(manager.task_aliases.lock().await.is_empty());
|
||||
assert_eq!(manager.completed_heals.lock().await.len(), 1);
|
||||
assert_eq!(manager.get_statistics().await.failed_tasks, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn cleanup_panic_is_supervised() {
|
||||
clear_scheduler_panic();
|
||||
let notice_bucket = "cleanup-panic-mrf";
|
||||
let notice_object = "object";
|
||||
let _ = rustfs_common::mrf_channel::take_mrf_repaired_events_for(notice_bucket);
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let manager = HealManager::new(storage, None);
|
||||
let mut request = HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal);
|
||||
request.source = HealRequestSource::Admin;
|
||||
let task_id = request.id.clone();
|
||||
assert_eq!(
|
||||
manager
|
||||
.submit_heal_request(request)
|
||||
.await
|
||||
.expect("cleanup request should be admitted"),
|
||||
HealAdmissionResult::Accepted
|
||||
);
|
||||
manager
|
||||
.mrf_repair_notice_targets
|
||||
.lock()
|
||||
.expect("mrf repair notice registry poisoned")
|
||||
.insert(
|
||||
task_id.clone(),
|
||||
vec![MrfRepairNoticeTarget {
|
||||
bucket: Arc::from(notice_bucket),
|
||||
object: Arc::from(notice_object),
|
||||
version_id: None,
|
||||
}],
|
||||
);
|
||||
arm_scheduler_panic(SchedulerPanicPoint::Cleanup, &task_id);
|
||||
process_manager_queue_once(&manager).await;
|
||||
|
||||
let status = tokio::time::timeout(Duration::from_secs(1), async {
|
||||
loop {
|
||||
if let Ok(status) = manager.get_task_status(&task_id).await
|
||||
&& matches!(status, HealTaskStatus::Completed)
|
||||
{
|
||||
break status;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("cleanup panic should leave a terminal status");
|
||||
clear_scheduler_panic();
|
||||
assert_eq!(status, HealTaskStatus::Completed);
|
||||
assert_eq!(manager.get_active_task_count().await, 0);
|
||||
assert!(manager.task_aliases.lock().await.is_empty());
|
||||
assert_eq!(manager.completed_heals.lock().await.len(), 1);
|
||||
assert_eq!(manager.get_statistics().await.successful_tasks, 1);
|
||||
let events = rustfs_common::mrf_channel::take_mrf_repaired_events_for(notice_bucket);
|
||||
assert_eq!(events.len(), 1, "cleanup panic must preserve successful MRF notice delivery");
|
||||
assert_eq!(events[0].object.as_ref(), notice_object);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn cancelled_retry_child_panic_does_not_rearchive_failed_status() {
|
||||
let manager = HealManager::new(Arc::new(MockStorage), None);
|
||||
let request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None);
|
||||
let task_id = request.id.clone();
|
||||
let retry_cancel_token = insert_retrying_request(&manager, request.clone()).await;
|
||||
|
||||
manager
|
||||
.cancel_task(&task_id)
|
||||
.await
|
||||
.expect("retry cancellation should succeed");
|
||||
assert!(retry_cancel_token.is_cancelled());
|
||||
|
||||
let state = PanicCleanupState {
|
||||
active_heals: manager.active_heals.clone(),
|
||||
heal_queue: manager.heal_queue.clone(),
|
||||
completed_heals: manager.completed_heals.clone(),
|
||||
task_aliases: manager.task_aliases.clone(),
|
||||
retrying_heals: manager.retrying_heals.clone(),
|
||||
mrf_repair_notice_targets: manager.mrf_repair_notice_targets.clone(),
|
||||
replacement_recovery_anchors: manager.replacement_recovery_anchors.clone(),
|
||||
statistics: manager.statistics.clone(),
|
||||
};
|
||||
finish_panicked_retry_child(task_id.clone(), request.heal_type, retry_cancel_token, state).await;
|
||||
|
||||
assert!(manager.retrying_heals.lock().await.is_empty());
|
||||
assert!(manager.completed_heals.lock().await.is_empty());
|
||||
assert_eq!(manager.get_statistics().await.failed_tasks, 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn active_cancel_wins_parent_panic_cleanup_without_completed_status() {
|
||||
let manager = HealManager::new(Arc::new(MockStorage), None);
|
||||
let request = HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal);
|
||||
let task_id = request.id.clone();
|
||||
let task = Arc::new(HealTask::from_request(request, Arc::new(MockStorage)));
|
||||
manager.active_heals.lock().await.insert(task_id.clone(), task.clone());
|
||||
|
||||
manager
|
||||
.cancel_task(&task_id)
|
||||
.await
|
||||
.expect("active task cancellation should win");
|
||||
assert_eq!(task.get_status().await, HealTaskStatus::Cancelled);
|
||||
|
||||
let state = PanicCleanupState {
|
||||
active_heals: manager.active_heals.clone(),
|
||||
heal_queue: manager.heal_queue.clone(),
|
||||
completed_heals: manager.completed_heals.clone(),
|
||||
task_aliases: manager.task_aliases.clone(),
|
||||
retrying_heals: manager.retrying_heals.clone(),
|
||||
mrf_repair_notice_targets: manager.mrf_repair_notice_targets.clone(),
|
||||
replacement_recovery_anchors: manager.replacement_recovery_anchors.clone(),
|
||||
statistics: manager.statistics.clone(),
|
||||
};
|
||||
finish_panicked_heal_task(task, task_id, state).await;
|
||||
|
||||
assert!(manager.completed_heals.lock().await.is_empty());
|
||||
assert_eq!(manager.get_statistics().await.failed_tasks, 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_duplicate_admission_is_atomic_with_queue_to_active_transition() {
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
|
||||
Reference in New Issue
Block a user