From ae87ddbe2f49e0387e06982cf1184fcfd7d9c1a4 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sun, 6 Sep 2026 04:26:28 +0800 Subject: [PATCH] test(rpc): reproduce same-UUID cross-instance rename --- rustfs/src/storage/rpc/node_service.rs | 222 +++++++++++++++++++++++++ 1 file changed, 222 insertions(+) diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 7278991c5..a7e0df269 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -4530,6 +4530,228 @@ mod tests { assert!(rename_response.error.is_some()); } + #[tokio::test] + async fn rename_data_same_uuid_uses_captured_instance_instead_of_global_disk() { + use crate::storage::storage_api::{ + ECStore, + ecstore_disk::{DiskAPI, RUSTFS_META_BUCKET, ReadOptions}, + init_local_disks_with_instance_ctx, new_instance_ctx, + }; + use rustfs_filemeta::{FileInfo, ObjectPartInfo}; + use tokio_util::sync::CancellationToken; + + async fn build_store(root: &std::path::Path) -> Arc { + let mut endpoints = Vec::new(); + for index in 0..4 { + let path = root.join(format!("disk{index}")); + tokio::fs::create_dir_all(&path).await.expect("create instance disk"); + let mut endpoint = Endpoint::try_from(path.to_str().expect("UTF-8 disk path")).expect("local endpoint"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(index); + endpoints.push(endpoint); + } + let pools = EndpointServerPools(vec![PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 4, + endpoints: Endpoints::from(endpoints), + cmd_line: "namespace-target-context".to_string(), + platform: "test".to_string(), + }]); + let instance = new_instance_ctx(); + init_local_disks_with_instance_ctx(&instance, pools.clone()) + .await + .expect("register this instance's real disks"); + // Match the isolated ECStore fixtures: startup still runs, while + // unrelated background recovery is cancelled for this process. + let shutdown = CancellationToken::new(); + shutdown.cancel(); + ECStore::new_with_instance_ctx("127.0.0.1:0".parse().expect("local address"), pools, shutdown, instance) + .await + .expect("initialize isolated ECStore") + } + + async fn context(store: &Arc) -> Arc { + ObjectStore::new(store.clone()) + .save_iam_config(serde_json::json!({"version": 1}), format!("{}/format.json", *IAM_CONFIG_PREFIX)) + .await + .expect("seed isolated IAM format"); + let iam = rustfs_iam::build_iam_sys(store.clone()).await.expect("build isolated IAM"); + Arc::new(crate::runtime_sources::AppContext::with_default_interfaces( + store.clone(), + iam, + Arc::new(KmsServiceManager::new()), + )) + } + + fn file_info(object: &str, version: Uuid, body: Bytes) -> FileInfo { + let mut fi = FileInfo::new(object, 1, 0); + fi.erasure.index = 1; + fi.name = object.to_string(); + fi.version_id = Some(version); + fi.size = i64::try_from(body.len()).expect("small fixture body"); + fi.parts = vec![ObjectPartInfo { + number: 1, + size: body.len(), + actual_size: fi.size, + ..Default::default() + }]; + fi.data = Some(body); + fi.set_inline_data(); + fi.mod_time = Some(OffsetDateTime::now_utc()); + fi + } + + // Both global publications are first-writer-wins. Run this fixture in + // its own nextest process; do not reset or replace another test's state. + assert!( + crate::runtime_sources::current_app_context().is_none(), + "requires an unpublished AppContext" + ); + super::timeout(Duration::from_secs(90), async { + let root_b = tempfile::tempdir().expect("instance B directory"); + let root_a = tempfile::tempdir().expect("instance A directory"); + let store_b = build_store(root_b.path()).await; + let context_b = context(&store_b).await; + let published = crate::runtime_sources::publish_test_app_context(context_b.clone()); + assert!(Arc::ptr_eq(&published, &context_b), "B must win the process AppContext publication"); + + // Copy the actual initialized format set, without editing IDs. The + // request UUID must identify an active physical disk in BOTH stores. + for index in 0..4 { + let relative = format!("disk{index}/{RUSTFS_META_BUCKET}/format.json"); + let target = root_a.path().join(&relative); + tokio::fs::create_dir_all(target.parent().expect("format parent")) + .await + .expect("create A format directory"); + tokio::fs::copy(root_b.path().join(&relative), target) + .await + .expect("copy the real disk format to A"); + } + let store_a = build_store(root_a.path()).await; + let service = make_server_for_context(Some(context(&store_a).await)); + assert!(Arc::ptr_eq(&service.resolve_object_store().expect("captured store"), &store_a)); + assert!(Arc::ptr_eq( + &crate::runtime_sources::current_object_store_handle().expect("global store"), + &store_b + )); + let disk_a = store_a.disk_map[&0][0].as_ref().expect("A disk zero").clone(); + let disk_b = store_b.disk_map[&0][0].as_ref().expect("B disk zero").clone(); + assert!(disk_a.is_local() && disk_b.is_local()); + assert!(!Arc::ptr_eq(&disk_a, &disk_b)); + let disk_id = disk_a.get_disk_id().await.expect("A disk ID").expect("formatted A disk"); + assert!(!disk_id.is_nil()); + assert_eq!(disk_b.get_disk_id().await.expect("B disk ID"), Some(disk_id)); + let global_disk = super::find_local_disk_by_ref(&disk_id.to_string()) + .await + .expect("global UUID lookup must resolve B before the request"); + assert!(Arc::ptr_eq(&global_disk, &disk_b)); + + let volume = "namespace-target-context"; + let object = "destination"; + let staging = "staged"; + let version = Uuid::new_v4(); + let new_body = Bytes::from_static(b"committed-through-captured-A"); + let new_fi = file_info(object, version, new_body.clone()); + let opts = ReadOptions { read_data: true, ..Default::default() }; + for (disk, old_body) in [ + (&disk_a, Bytes::from_static(b"old-body-A")), + (&disk_b, Bytes::from_static(b"old-body-B")), + ] { + disk.make_volume(volume).await.expect("create destination volume"); + disk.write_metadata(volume, volume, object, file_info(object, version, old_body.clone())) + .await + .expect("write real old object metadata and inline body"); + disk.write_metadata(volume, volume, staging, new_fi.clone()) + .await + .expect("stage identical valid metadata on both physical disks"); + let seeded = disk + .read_version(volume, volume, object, &version.to_string(), &opts) + .await + .expect("decode seeded inline object before invoking the handler"); + assert_eq!(seeded.data, Some(old_body), "the real reader must return the seeded body"); + } + let a_meta = disk_a.path().join(volume).join(object).join("xl.meta"); + let b_meta = disk_b.path().join(volume).join(object).join("xl.meta"); + let a_staging = disk_a.path().join(volume).join(staging).join("xl.meta"); + let b_staging = disk_b.path().join(volume).join(staging).join("xl.meta"); + let a_before = tokio::fs::read(&a_meta).await.expect("A old metadata bytes"); + let b_before = tokio::fs::read(&b_meta).await.expect("B old metadata bytes"); + let b_staging_before = tokio::fs::read(&b_staging).await.expect("B staged metadata bytes"); + assert!(tokio::fs::try_exists(&a_staging).await.expect("A staging exists")); + assert_ne!(a_before, b_before, "the old on-disk bodies must distinguish A from B"); + super::timeout(Duration::from_secs(10), async { + while store_a.scanner_data_usage_publication_blocked().await + || store_b.scanner_data_usage_publication_blocked().await + { + tokio::task::yield_now().await; + } + }) + .await + .expect("startup namespace commits must drain before measuring the handler"); + let generation_before = ( + store_a.scanner_namespace_mutation_generation(), + store_b.scanner_namespace_mutation_generation(), + ); + + let mut request = Request::new(RenameDataRequest { + disk: disk_id.to_string(), + src_volume: volume.to_string(), + src_path: staging.to_string(), + dst_volume: volume.to_string(), + dst_path: object.to_string(), + file_info: serde_json::to_string(&new_fi).expect("encode real FileInfo"), + file_info_bin: Vec::new().into(), + scanner_publication_lease_token: Vec::new().into(), + }); + let body = rustfs_protos::canonical_rename_data_request_body(request.get_ref()).expect("canonical rename body"); + set_tonic_canonical_body_digest(&mut request, &body).expect("body-bound handler request"); + let response = super::timeout(Duration::from_secs(10), service.rename_data(request)) + .await + .expect("real rename handler must finish within ten seconds") + .expect("rename handler response") + .into_inner(); + assert!(response.success, "the valid staged rename must execute: {:?}", response.error); + + let a_after = disk_a + .read_version(volume, volume, object, &version.to_string(), &opts) + .await + .expect("read A's physical object after rename"); + let b_after = disk_b + .read_version(volume, volume, object, &version.to_string(), &opts) + .await + .expect("read B's physical object after rename"); + let a_bytes_after = tokio::fs::read(&a_meta).await.expect("A metadata after rename"); + let b_bytes_after = tokio::fs::read(&b_meta).await.expect("B metadata after rename"); + let b_staging_after = tokio::fs::read(&b_staging).await.ok(); + let generation_after = ( + store_a.scanner_namespace_mutation_generation(), + store_b.scanner_namespace_mutation_generation(), + ); + let pending_after = ( + store_a.scanner_data_usage_publication_blocked().await, + store_b.scanner_data_usage_publication_blocked().await, + ); + for disk in store_a.disk_map.values().chain(store_b.disk_map.values()).flatten().flatten() { + disk.close().await.expect("close real fixture disk before assertions and directory cleanup"); + } + assert_eq!( + a_after.data, + Some(new_body), + "RenameData must commit to captured A, even when global B owns the same UUID; B body={:?}, generations={generation_before:?}->{generation_after:?}, pending={pending_after:?}", + b_after.data + ); + assert_ne!(a_bytes_after, a_before, "A metadata must actually be replaced"); + assert_eq!(b_after.data, Some(Bytes::from_static(b"old-body-B")), "B body must remain unchanged"); + assert_eq!(b_bytes_after, b_before, "B metadata must remain byte-for-byte unchanged"); + assert_eq!(b_staging_after, Some(b_staging_before), "B staging must not be consumed"); + assert_eq!(pending_after, (false, false), "both stores must reach a stable terminal state"); + }) + .await + .expect("two-instance handler fixture must finish within ninety seconds"); + } + #[tokio::test] async fn test_make_volumes_invalid_disk() { let service = create_test_node_service();