From 5cedf73d7c31fc48db156440df30e63d6fa8a710 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sun, 23 Aug 2026 21:32:39 +0800 Subject: [PATCH] fix(ecstore): terminate walk directory streams (#6462) * fix(ecstore): terminate walk directory streams * test(e2e): refresh node service selection * fix(filemeta): initialize empty metacache streams --- .config/e2e-full-selection.txt | 4 +- .config/nextest.toml | 4 +- .../src/reliant/node_interact_test.rs | 327 ++++++++++-------- crates/ecstore/src/disk/local.rs | 1 + crates/filemeta/src/metacache.rs | 11 + docs/testing/e2e-suite-inventory.md | 13 +- 6 files changed, 197 insertions(+), 163 deletions(-) diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index 59a5d8ab9..97cc7e810 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=f832043fcca8c0b616c5d820a3a652da7544298ef5812a8668a3a9a3e4607b8b -sha256-linux=93b94adb110b86a41d0b7313909e0bf53cb1515e2d08e8f105652b29b249990f +sha256-darwin=b63c54946d978e22660762e8e0351f6839bebf9be2282fe08e03e775e4303f0a +sha256-linux=ef3be856bd3257c2c369428f66a48ad073dad8187229fe95840a600a81edf22b diff --git a/.config/nextest.toml b/.config/nextest.toml index 841368bef..d664e731c 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -315,7 +315,7 @@ slow-timeout = { period = "60s", terminate-after = 2, grace-period = "10s" } # the target incl. multipart and the resync path, SSE-C and # target-without-KMS stay fail-closed), and one guards event/history # observers. -# * 12 `_real_dual_node` site-replication tests — each spawns TWO full rustfs +# * 13 `_real_dual_node` site-replication tests — each spawns TWO full rustfs # servers and drives the cross-process site-replication control plane. # * 1 `_real_three_node` site-replication test. # * 1 `_real_single_node` service-account round-trip test. @@ -406,7 +406,7 @@ path = "junit.xml" # object_lambda) — too heavy for the merge budget; they run in the # e2e-nightly serial cluster-fault lane. # * replication_extension_test — repl-1 already splits it into the PR -# `e2e-smoke` (20 fast) and `e2e-repl-nightly` (55 slow) lanes and reserves +# `e2e-smoke` (20 fast) and `e2e-repl-nightly` (56 slow) lanes and reserves # it for those, so e2e-full does not double-run it. # * #[ignore]d tests — nextest skips them by default (no --run-ignored); the # manual-localhost:9000 reliant/policy tests are ci-13's migration. diff --git a/crates/e2e_test/src/reliant/node_interact_test.rs b/crates/e2e_test/src/reliant/node_interact_test.rs index e5273bc71..3ded227d4 100644 --- a/crates/e2e_test/src/reliant/node_interact_test.rs +++ b/crates/e2e_test/src/reliant/node_interact_test.rs @@ -13,207 +13,228 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::common::workspace_root; +use crate::common::RustFSTestEnvironment; use crate::storage_api::node_interact::{ TonicInterceptor, VolumeInfo, WalkDirOptions, gen_tonic_signature_interceptor, node_service_time_out_client, }; -use futures::future::join_all; +use aws_sdk_s3::primitives::ByteStream; use rmp_serde::{Deserializer, Serializer}; -use rustfs_filemeta::{MetaCacheEntry, MetacacheReader, MetacacheWriter}; +use rustfs_filemeta::MetaCacheEntry; use rustfs_protos::proto_gen::node_service::WalkDirRequest; use rustfs_protos::{ models::{PingBody, PingBodyBuilder}, - proto_gen::node_service::{ - ListVolumesRequest, LocalStorageInfoRequest, MakeVolumeRequest, PingRequest, PingResponse, ReadAllRequest, - }, + proto_gen::node_service::{ListVolumesRequest, LocalStorageInfoRequest, MakeVolumeRequest, PingRequest, ReadAllRequest}, }; use serde::{Deserialize, Serialize}; use std::error::Error; use std::io::Cursor; -use std::path::PathBuf; -use tokio::spawn; use tonic::Request; use tonic::codegen::tokio_stream::StreamExt; -const CLUSTER_ADDR: &str = "http://localhost:9000"; +type TestResult = Result<(), Box>; + +const TEST_RPC_SECRET: &str = "rustfs-internode-signature-e2e-secret"; fn signature_interceptor() -> TonicInterceptor { TonicInterceptor::Signature(gen_tonic_signature_interceptor()) } +fn rpc_client_error(error: Box) -> std::io::Error { + std::io::Error::other(error.to_string()) +} + +async fn start_server() -> Result> { + let _ = rustfs_credentials::set_global_rpc_secret(TEST_RPC_SECRET.to_string()); + let effective = rustfs_credentials::try_get_rpc_token().expect("RPC secret must resolve in the test process"); + assert_eq!(effective, TEST_RPC_SECRET, "the test process uses an unexpected RPC secret"); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server_without_cleanup_with_env(&[ + ("RUSTFS_RPC_SECRET", TEST_RPC_SECRET), + ("RUSTFS_INTERNODE_RPC_SIGNATURE_STRICT", "false"), + ("RUSTFS_INTERNODE_RPC_BODY_DIGEST_STRICT", "false"), + ("RUSTFS_INTERNODE_RPC_REPLAY_SCOPE_STRICT", "false"), + ("RUST_LOG", "error"), + ]) + .await?; + Ok(env) +} + #[tokio::test] -#[ignore = "requires running RustFS server at localhost:9000"] -async fn ping() -> Result<(), Box> { +async fn ping() -> TestResult { + let env = start_server().await?; let mut fbb = flatbuffers::FlatBufferBuilder::new(); let payload = fbb.create_vector(b"hello world"); - let mut builder = PingBodyBuilder::new(&mut fbb); builder.add_payload(payload); let root = builder.finish(); fbb.finish(root, None); - let finished_data = fbb.finished_data(); - - let decoded_payload = flatbuffers::root::(finished_data); - assert!(decoded_payload.is_ok()); - - // Create client - let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), signature_interceptor()).await?; - - // Construct PingRequest - let request = Request::new(PingRequest { - version: 1, - body: bytes::Bytes::copy_from_slice(finished_data), - }); - - // Send request and get response - let response: PingResponse = client.ping(request).await?.into_inner(); - - // Print response - let ping_response_body = flatbuffers::root::(&response.body); - if let Err(e) = ping_response_body { - eprintln!("{e}"); - } else { - println!("ping_resp:body(flatbuffer): {ping_response_body:?}"); - } + let mut client = node_service_time_out_client(&env.url, signature_interceptor()) + .await + .map_err(rpc_client_error)?; + let response = client + .ping(Request::new(PingRequest { + version: 1, + body: bytes::Bytes::copy_from_slice(fbb.finished_data()), + })) + .await? + .into_inner(); + assert_eq!(response.version, 1); + let body = flatbuffers::root::(&response.body)?; + assert_eq!(body.payload().expect("ping response must contain a payload").bytes(), b"hello, caller"); Ok(()) } #[tokio::test] -#[ignore = "requires running RustFS server at localhost:9000"] -async fn make_volume() -> Result<(), Box> { - let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), signature_interceptor()).await?; - let request = Request::new(MakeVolumeRequest { - disk: "data".to_string(), - volume: "dandan".to_string(), - }); +async fn make_volume() -> TestResult { + let env = start_server().await?; + let mut client = node_service_time_out_client(&env.url, signature_interceptor()) + .await + .map_err(rpc_client_error)?; + let response = client + .make_volume(Request::new(MakeVolumeRequest { + disk: env.temp_dir.clone(), + volume: "node-rpc-volume".to_string(), + })) + .await? + .into_inner(); - let response = client.make_volume(request).await?.into_inner(); - if response.success { - println!("success"); - } else { - println!("failed: {:?}", response.error); - } + assert!(response.success, "make_volume failed: {:?}", response.error); + assert!(std::path::Path::new(&env.temp_dir).join("node-rpc-volume").is_dir()); Ok(()) } #[tokio::test] -#[ignore = "requires running RustFS server at localhost:9000"] -async fn list_volumes() -> Result<(), Box> { - let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), signature_interceptor()).await?; - let request = Request::new(ListVolumesRequest { - disk: "data".to_string(), - }); +async fn list_volumes() -> TestResult { + let env = start_server().await?; + let mut client = node_service_time_out_client(&env.url, signature_interceptor()) + .await + .map_err(rpc_client_error)?; + let created = client + .make_volume(Request::new(MakeVolumeRequest { + disk: env.temp_dir.clone(), + volume: "node-rpc-listed-volume".to_string(), + })) + .await? + .into_inner(); + assert!(created.success, "make_volume failed: {:?}", created.error); - let response = client.list_volumes(request).await?.into_inner(); - let volume_infos: Vec = response + let response = client + .list_volumes(Request::new(ListVolumesRequest { + disk: env.temp_dir.clone(), + })) + .await? + .into_inner(); + assert!(response.success, "list_volumes failed: {:?}", response.error); + let volumes = response .volume_infos - .into_iter() - .filter_map(|json_str| serde_json::from_str::(&json_str).ok()) - .collect(); - - println!("{volume_infos:?}"); + .iter() + .map(|json| serde_json::from_str::(json)) + .collect::, _>>()?; + assert!(volumes.iter().any(|volume| volume.name == "node-rpc-listed-volume")); Ok(()) } #[tokio::test] -#[ignore = "requires running RustFS server at localhost:9000"] -async fn walk_dir() -> Result<(), Box> { - println!("walk_dir"); - // TODO: use writer +async fn walk_dir() -> TestResult { + let env = start_server().await?; + let s3 = env.create_s3_client(); + let bucket = "node-rpc-walk-bucket"; + let key = "prefix/object.txt"; + env.create_test_bucket(bucket).await?; + s3.put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"walk payload")) + .send() + .await?; + let opts = WalkDirOptions { - bucket: "dandan".to_owned(), - base_dir: "".to_owned(), + bucket: bucket.to_string(), recursive: true, ..Default::default() }; - let (rd, mut wr) = tokio::io::duplex(1024); - let mut buf = Vec::new(); - opts.serialize(&mut Serializer::new(&mut buf))?; - let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), signature_interceptor()).await?; - let disk_path = std::env::var_os("RUSTFS_DISK_PATH").map(PathBuf::from).unwrap_or_else(|| { - let mut path = workspace_root(); - path.push("target"); - path.push(if cfg!(debug_assertions) { "debug" } else { "release" }); - path.push("data"); - path - }); - let request = Request::new(WalkDirRequest { - disk: disk_path.to_string_lossy().into_owned(), - walk_dir_options: buf.into(), - }); - let mut response = client.walk_dir(request).await?.into_inner(); + let mut encoded = Vec::new(); + opts.serialize(&mut Serializer::new(&mut encoded))?; + let mut client = node_service_time_out_client(&env.url, signature_interceptor()) + .await + .map_err(rpc_client_error)?; + let mut stream = client + .walk_dir(Request::new(WalkDirRequest { + disk: env.temp_dir.clone(), + walk_dir_options: encoded.into(), + })) + .await? + .into_inner(); - let job1 = spawn(async move { - let mut out = MetacacheWriter::new(&mut wr); - loop { - match response.next().await { - Some(Ok(resp)) => { - if !resp.success { - println!("{}", resp.error_info.unwrap_or_else(|| "".to_string())); - } - let entry = serde_json::from_str::(&resp.meta_cache_entry) - .map_err(|_e| std::io::Error::other(format!("Unexpected response: {response:?}"))) - .unwrap(); - out.write_obj(&entry).await.unwrap(); - } - None => { - let _ = out.close().await; - break; - } - _ => { - println!("Unexpected response: {response:?}"); - let _ = out.close().await; - break; - } - } - } - }); - let job2 = spawn(async move { - let mut reader = MetacacheReader::new(rd); - while let Ok(Some(entry)) = reader.peek().await { - println!("{entry:?}"); - } - }); - - join_all(vec![job1, job2]).await; - Ok(()) -} - -#[tokio::test] -#[ignore = "requires running RustFS server at localhost:9000"] -async fn read_all() -> Result<(), Box> { - let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), signature_interceptor()).await?; - let request = Request::new(ReadAllRequest { - disk: "data".to_string(), - volume: "ff".to_string(), - path: "format.json".to_string(), - }); - - let response = client.read_all(request).await?.into_inner(); - let volume_infos = response.data; - - println!("{}", response.success); - println!("{volume_infos:?}"); - Ok(()) -} - -#[tokio::test] -#[ignore = "requires running RustFS server at localhost:9000"] -async fn storage_info() -> Result<(), Box> { - let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), signature_interceptor()).await?; - let request = Request::new(LocalStorageInfoRequest { metrics: true }); - - let response = client.local_storage_info(request).await?.into_inner(); - if !response.success { - println!("{:?}", response.error_info); - return Ok(()); + let mut entries = Vec::new(); + while let Some(response) = stream.next().await { + let response = response?; + assert!(response.success, "walk_dir failed: {:?}", response.error_info); + entries.push(serde_json::from_str::(&response.meta_cache_entry)?); } - let info = response.storage_info; - - let mut buf = Deserializer::new(Cursor::new(info)); - let storage_info: rustfs_madmin::StorageInfo = Deserialize::deserialize(&mut buf).unwrap(); - println!("{storage_info:?}"); + assert!( + entries.iter().any(|entry| entry.name == key), + "walk_dir did not return {key}: {entries:?}" + ); + Ok(()) +} + +#[tokio::test] +async fn read_all() -> TestResult { + let env = start_server().await?; + let mut client = node_service_time_out_client(&env.url, signature_interceptor()) + .await + .map_err(rpc_client_error)?; + let volume = "node-rpc-read-volume"; + let created = client + .make_volume(Request::new(MakeVolumeRequest { + disk: env.temp_dir.clone(), + volume: volume.to_string(), + })) + .await? + .into_inner(); + assert!(created.success, "make_volume failed: {:?}", created.error); + tokio::fs::write(std::path::Path::new(&env.temp_dir).join(volume).join("payload.bin"), b"read payload").await?; + + let response = client + .read_all(Request::new(ReadAllRequest { + disk: env.temp_dir.clone(), + volume: volume.to_string(), + path: "payload.bin".to_string(), + })) + .await? + .into_inner(); + assert!(response.success, "read_all failed: {:?}", response.error); + assert_eq!(response.data.as_ref(), b"read payload"); + Ok(()) +} + +#[tokio::test] +async fn storage_info() -> TestResult { + let env = start_server().await?; + let mut client = node_service_time_out_client(&env.url, signature_interceptor()) + .await + .map_err(rpc_client_error)?; + let response = client + .local_storage_info(Request::new(LocalStorageInfoRequest { metrics: true })) + .await? + .into_inner(); + assert!(response.success, "local_storage_info failed: {:?}", response.error_info); + + let mut decoder = Deserializer::new(Cursor::new(response.storage_info)); + let storage_info: rustfs_madmin::StorageInfo = Deserialize::deserialize(&mut decoder)?; + let expected_disk = std::fs::canonicalize(&env.temp_dir)?; + assert!(!storage_info.disks.is_empty(), "local_storage_info returned no disks"); + assert!( + storage_info + .disks + .iter() + .any(|disk| std::path::Path::new(&disk.drive_path) == expected_disk), + "local_storage_info did not include the configured disk: {:?}", + storage_info.disks + ); Ok(()) } diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 5217e819b..0290a040d 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -8963,6 +8963,7 @@ impl DiskAPI for LocalDisk { ) .await?; + out.close().await?; Ok(()) } diff --git a/crates/filemeta/src/metacache.rs b/crates/filemeta/src/metacache.rs index 427be0027..0571bc466 100644 --- a/crates/filemeta/src/metacache.rs +++ b/crates/filemeta/src/metacache.rs @@ -894,6 +894,7 @@ impl MetacacheWriter { } pub async fn close(&mut self) -> Result<()> { + self.init().await?; rmp::encode::write_bool(&mut self.buf, false).map_err(|e| Error::other(format!("{e:?}")))?; self.flush().await?; Ok(()) @@ -1285,6 +1286,16 @@ mod tests { assert_eq!(objs, nobjs); } + #[tokio::test] + async fn empty_writer_emits_a_valid_stream() { + let mut output = Cursor::new(Vec::new()); + let mut writer = MetacacheWriter::new(&mut output); + writer.close().await.expect("empty stream should close"); + + let mut reader = MetacacheReader::new(Cursor::new(output.into_inner())); + assert!(reader.read_all().await.expect("empty stream should decode").is_empty()); + } + fn corrupt_stream_with_metadata_len(len_marker: u8, len: u32) -> Vec { let mut data = Vec::new(); rmp::encode::write_u8(&mut data, METACACHE_STREAM_VERSION).unwrap(); diff --git a/docs/testing/e2e-suite-inventory.md b/docs/testing/e2e-suite-inventory.md index 3ede11cee..430a43d12 100644 --- a/docs/testing/e2e-suite-inventory.md +++ b/docs/testing/e2e-suite-inventory.md @@ -28,9 +28,9 @@ | bucket_stats_regression_test | 3 | | | chaos | 2 | | | checksum_upload_test | 7 | | -| cluster_concurrency_test | 2 | 🌙 | +| cluster_concurrency_test | 3 | 🌙 | | cluster_multidrive_pool_test | 2 | 🌙 | -| common | 14 | | +| common | 16 | | | compression_test | 6 | ✅ | | connection_cap_test | 2 | | | console_smoke_test | 1 | ✅ | @@ -58,7 +58,7 @@ | heal_erasure_disk_rebuild_test | 4 | 🌙 | | inline_fast_path_cluster_test | 16 | | | internode_rpc_signature_e2e_test | 5 | | -| kms | 46 | | +| kms | 48 | | | leading_slash_key_test | 2 | ✅ | | lifecycle_regression_test | 4 | | | list_buckets_auth_test | 1 | ✅ | @@ -84,8 +84,9 @@ | protocols | 16 | 🌙 | | quota_test | 14 | | | reliability_disk_fault_test | 4 | | -| reliant | 25 | 19 ✅ | -| replication_extension_test | 75 | 20 ✅ +55 🌙 | +| reliant | 31 | 19 ✅ | +| replication_extension_test | 76 | 20 ✅ +56 🌙 | +| replication_lww_receiver_test | 1 | | | security_boundary_test | 4 | | | server_startup_failfast_test | 1 | | | snowball_auto_extract_test | 6 | | @@ -99,4 +100,4 @@ | tls_hot_reload_test | 1 | ✅ | | version_id_regression_test | 10 | ✅ | -**Total listed: 578 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 456 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23. +**Total listed: 591 tests across 83 modules · PR smoke: 163 tests / 36 modules · merge/main full: 467 tests / 74 modules · nightly replication: 56 tests · nightly cluster faults: 29 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.