diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 5eed710c6..206218492 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -1824,7 +1824,7 @@ mod tests { let pool_meta = PoolMeta::new(&store.pools, &PoolMeta::default()); pool_meta - .save(store.pools.clone()) + .save_for_startup(store.pools.clone()) .await .expect("pool metadata should be persisted before format heal"); *store.pool_meta_save_gate.lock().await = PoolMetaWriteState::default(); diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index 4dd777e15..4521e087f 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -937,6 +937,7 @@ mod tests { POOL_META_IDENTITY_NAME, POOL_META_NAME, POOL_META_VERSION, PoolDecommissionInfo, PoolStatus, initialized_pool_meta_identity_for_test, pool_meta_v3_commit_state_for_test, }; + use crate::core::sets::Sets; use crate::disk::error::DiskError; use crate::layout::endpoint::Endpoint; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; @@ -2320,17 +2321,78 @@ mod tests { async fn persist_reload_snapshot(store: &ECStore, snapshot: &PoolMeta) { snapshot - .save(store.pools.clone()) + .save_for_startup(store.pools.clone()) .await .expect("pool meta snapshot should persist to every pool"); } + #[derive(Debug)] + struct PoolMetaCommitPauseStorage { + inner: Arc, + pause_committed_pool_meta: bool, + paused: Arc, + release: Arc, + } + + #[async_trait::async_trait] + impl crate::storage_api_contracts::object::ObjectIO for PoolMetaCommitPauseStorage { + type Error = Error; + type RangeSpec = crate::storage_api_contracts::range::HTTPRangeSpec; + type HeaderMap = http::HeaderMap; + type ObjectOptions = crate::object_api::ObjectOptions; + type ObjectInfo = crate::object_api::ObjectInfo; + type GetObjectReader = crate::object_api::GetObjectReader; + type PutObjectReader = crate::object_api::PutObjReader; + + async fn get_object_reader( + &self, + bucket: &str, + object: &str, + range: Option, + h: Self::HeaderMap, + opts: &Self::ObjectOptions, + ) -> std::result::Result { + self.inner.get_object_reader(bucket, object, range, h, opts).await + } + + async fn put_object( + &self, + bucket: &str, + object: &str, + data: &mut Self::PutObjectReader, + opts: &Self::ObjectOptions, + ) -> std::result::Result { + let result = self.inner.put_object(bucket, object, data, opts).await; + if result.is_ok() && object == POOL_META_NAME && self.pause_committed_pool_meta { + let committed = read_config_no_lock_preserve_empty_with_metadata(self.inner.clone(), POOL_META_NAME) + .await + .ok() + .and_then(|(payload, _)| pool_meta_v3_commit_state_for_test(payload).ok()) + .is_some_and(|(_, committed)| committed); + if committed { + self.paused.notify_one(); + self.release.notified().await; + } + } + result + } + } + #[tokio::test] #[serial_test::serial] async fn pool_meta_runtime_load_waits_for_multi_pool_commit_fence() { let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-load-fence", &[2, 2]).await; let old = PoolMeta::new(&store.pools, &PoolMeta::default()); - old.save(store.pools.clone()).await.expect("old pool metadata should persist"); + old.save_for_startup(store.pools.clone()) + .await + .expect("old V3 pool metadata should persist"); + let (old_payload, _) = read_config_no_lock_preserve_empty_with_metadata(store.pools[0].clone(), POOL_META_NAME) + .await + .expect("old V3 pool metadata should be readable"); + let (old_generation, old_committed) = + pool_meta_v3_commit_state_for_test(old_payload).expect("old replica should be a V3 envelope"); + assert!(old_committed, "old generation should be committed before the mixed-state save"); + let next_generation = old_generation.checked_add(1).expect("test generation should not exhaust"); let mut newer = old.clone(); newer.pools[0].last_update += TimeDuration::seconds(1); @@ -2342,11 +2404,6 @@ mod tests { .get_write_lock(get_lock_acquire_timeout()) .await .expect("pool metadata write fence should be acquired"); - newer - .save_for_startup(vec![store.pools[0].clone()]) - .await - .expect("first replica should enter the new generation"); - let started = Arc::new(tokio::sync::Notify::new()); let mut load_task = tokio::spawn({ let store = store.clone(); @@ -2357,17 +2414,57 @@ mod tests { } }); started.notified().await; + let save_paused = Arc::new(tokio::sync::Notify::new()); + let save_release = Arc::new(tokio::sync::Notify::new()); + let save_task = tokio::spawn({ + let newer = newer.clone(); + let pools = vec![ + Arc::new(PoolMetaCommitPauseStorage { + inner: store.pools[0].clone(), + pause_committed_pool_meta: true, + paused: save_paused.clone(), + release: save_release.clone(), + }), + Arc::new(PoolMetaCommitPauseStorage { + inner: store.pools[1].clone(), + pause_committed_pool_meta: false, + paused: save_paused.clone(), + release: save_release.clone(), + }), + ]; + async move { newer.save_for_startup(pools).await } + }); + tokio::time::timeout(std::time::Duration::from_secs(5), save_paused.notified()) + .await + .expect("the newer V3 generation should pause after the first replica commit"); assert!( tokio::time::timeout(std::time::Duration::from_millis(100), &mut load_task) .await .is_err(), - "runtime load must not observe a valid-old/valid-new intermediate state" + "runtime load must not observe a partially committed V3 generation while fenced" + ); + let (committed_payload, _) = read_config_no_lock_preserve_empty_with_metadata(store.pools[0].clone(), POOL_META_NAME) + .await + .expect("first replica should be readable while the runtime read is fenced"); + let (pending_payload, _) = read_config_no_lock_preserve_empty_with_metadata(store.pools[1].clone(), POOL_META_NAME) + .await + .expect("second replica should be readable while the runtime read is fenced"); + assert_eq!( + pool_meta_v3_commit_state_for_test(committed_payload).expect("first replica should be a V3 envelope"), + (next_generation, true), + "first replica should hold the committed new generation before the fence releases" + ); + assert_eq!( + pool_meta_v3_commit_state_for_test(pending_payload).expect("second replica should be a V3 envelope"), + (next_generation, false), + "second replica should still hold the pending new generation before the fence releases" ); - newer - .save_for_startup(vec![store.pools[1].clone()]) + save_release.notify_one(); + save_task .await - .expect("second replica should enter the new generation"); + .expect("newer V3 save task should not panic") + .expect("the newer generation should finish while the runtime read is fenced"); drop(pool_meta_guard); let loaded = tokio::time::timeout(std::time::Duration::from_secs(5), load_task) diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index bda5d23d6..bcd589c54 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -14707,8 +14707,10 @@ mod tests { .expect("multipart object must be placed in one source pool"); if upload_pool != 0 { for (source_disk, target_disk) in pool_disk_paths[upload_pool].iter().zip(&pool_disk_paths[0]) { - std::fs::rename(source_disk.join(&bucket).join(object), target_disk.join(&bucket).join(object)) - .expect("normalize the test object into the old pool"); + let source_object = source_disk.join(&bucket).join(object); + let target_bucket = target_disk.join(&bucket); + std::fs::create_dir_all(&target_bucket).expect("create normalized target bucket directory"); + std::fs::rename(source_object, target_bucket.join(object)).expect("normalize the test object into the old pool"); } } let source_pool = 0; diff --git a/rustfs/tests/connect_registration.rs b/rustfs/tests/connect_registration.rs index 078aaf4a1..b8a8161fa 100644 --- a/rustfs/tests/connect_registration.rs +++ b/rustfs/tests/connect_registration.rs @@ -419,6 +419,20 @@ fn stores(temp: &tempfile::TempDir) -> (IdentityStore, CredentialStore) { ) } +#[cfg(target_os = "linux")] +fn secure_tempdir() -> tempfile::TempDir { + use std::os::unix::fs::PermissionsExt as _; + + let home = fs::canonicalize(std::env::var_os("HOME").expect("test requires a protected home directory")) + .expect("test home directory must resolve without symlink components"); + fs::set_permissions(&home, fs::Permissions::from_mode(0o700)) + .expect("test home directory permissions must be restricted to the process owner"); + tempfile::Builder::new() + .prefix(".connect-registration-") + .tempdir_in(home) + .expect("temporary directory inside the protected home directory") +} + fn stage_next_identity(temp: &tempfile::TempDir) -> DeviceIdentity { let directory = temp.path().join("identity"); fs::create_dir_all(&directory).expect("identity directory"); @@ -902,7 +916,7 @@ async fn heartbeat_runtime_respects_rotation_retry_after_without_blocking_heartb let mut status = runtime.status(); wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await; - tokio::time::sleep(Duration::from_millis(350)).await; + wait_for_requests(&server, 4).await; runtime.shutdown().await; let paths = server.paths.lock().expect("paths lock"); @@ -925,7 +939,7 @@ async fn heartbeat_runtime_skips_only_valid_pending_reenrollment() { Reply::Json(StatusCode::SERVICE_UNAVAILABLE, json!({})), Reply::Json(StatusCode::OK, heartbeat_response("2026-08-25T01:02:03Z")), ], - true, + false, ) .await; config.endpoint = server.endpoint.clone(); @@ -1086,7 +1100,7 @@ async fn heartbeat_retries_with_the_new_credential_after_concurrent_rotation() { #[cfg(target_os = "linux")] #[tokio::test] async fn inventory_retries_with_the_new_credential_after_concurrent_rotation() { - let temp = tempfile::tempdir().expect("temp dir"); + let temp = secure_tempdir(); let pki = TestPki::new(); let (mut config, current, next) = runtime_config(&temp, &pki, "https://localhost/agent/", 1); let (rotated, rotated_stored) = rotation_response(&pki, &next, 0x23); @@ -1181,7 +1195,7 @@ async fn inventory_retries_with_the_new_credential_after_concurrent_rotation() { #[cfg(target_os = "linux")] #[tokio::test] async fn inventory_first_recovers_a_saved_reenrollment_before_telemetry() { - let temp = tempfile::tempdir().expect("temp dir"); + let temp = secure_tempdir(); let pki = TestPki::new(); let failed = server(&pki, vec![Reply::Json(StatusCode::SERVICE_UNAVAILABLE, json!({})); 3]).await; let (mut config, _, next) = runtime_config(&temp, &pki, &failed.endpoint, 1);