fix(ecstore): replace panic-driven pool and set stubs (#3753)

* fix(ecstore): replace panic-driven pool and set stubs

* test(runtime): tolerate restricted local bind checks

* fix(ecstore): remove remaining trait stub placeholders

* fix(ecstore): tighten trait stub follow-up semantics

* chore: ignore local worktrees

* chore: update layer dependency baseline for resolve_* context entries

Add 7 accepted infra->app dependency entries introduced by recent
refactoring PRs (#3770, #3771, #3772) that route global state lookups
through app::context::resolve_* functions.

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-06-23 12:31:17 +08:00
committed by GitHub
parent c421e73fef
commit 583a23bdf2
20 changed files with 2361 additions and 162 deletions
+1
View File
@@ -66,3 +66,4 @@ tmp/
crates/*/docs
fuzz/target
outputs
worktrees/*
+14 -5
View File
@@ -765,9 +765,16 @@ mod tests {
net::TcpListener,
};
async fn capture_delete_objects_sha256_header() -> (String, tokio::task::JoinHandle<String>) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let endpoint = listener.local_addr().unwrap().to_string();
async fn capture_delete_objects_sha256_header() -> Option<(String, tokio::task::JoinHandle<String>)> {
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None,
Err(err) => panic!("test listener should bind: {err}"),
};
let endpoint = listener
.local_addr()
.expect("listener local address should be available")
.to_string();
let task = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.unwrap();
let mut request = Vec::new();
@@ -801,7 +808,7 @@ mod tests {
sha256_header
});
(endpoint, task)
Some((endpoint, task))
}
#[tokio::test]
@@ -813,7 +820,9 @@ mod tests {
}];
let body = generate_remove_multi_objects_request(&objects);
let expected = rustfs_utils::hex(HashAlgorithm::SHA256.hash_encode(&body));
let (endpoint, header_task) = capture_delete_objects_sha256_header().await;
let Some((endpoint, header_task)) = capture_delete_objects_sha256_header().await else {
return;
};
let client = TransitionClient::new(
&endpoint,
Options {
+1
View File
@@ -31,6 +31,7 @@ pub use internode_data_transport::build_internode_data_transport_from_env;
pub use peer_rest_client::{
PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC,
};
pub(crate) use peer_s3_client::heal_bucket_local_on_disks;
pub use peer_s3_client::{LocalPeerS3Client, PeerS3Client, S3PeerSys};
pub use remote_disk::RemoteDisk;
pub use remote_locker::RemoteClient;
+5 -1
View File
@@ -954,7 +954,11 @@ pub async fn heal_bucket_local(bucket: &str, opts: &HealOpts) -> Result<HealResu
heal_bucket_local_on_disks(bucket, opts, disks).await
}
async fn heal_bucket_local_on_disks(bucket: &str, opts: &HealOpts, disks: Vec<Option<DiskStore>>) -> Result<HealResultItem> {
pub(crate) async fn heal_bucket_local_on_disks(
bucket: &str,
opts: &HealOpts,
disks: Vec<Option<DiskStore>>,
) -> Result<HealResultItem> {
let before_state = Arc::new(RwLock::new(vec![String::new(); disks.len()]));
let after_state = Arc::new(RwLock::new(vec![String::new(); disks.len()]));
+18 -6
View File
@@ -2662,8 +2662,12 @@ mod tests {
#[tokio::test]
async fn test_remote_disk_is_online_detects_active_listener() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
let url = url::Url::parse(&format!("http://{}:{}/data/rustfs0", addr.ip(), addr.port())).unwrap();
let endpoint = Endpoint {
@@ -2691,8 +2695,12 @@ mod tests {
async fn test_remote_disk_is_online_detects_missing_listener() {
init_tracing(Level::ERROR);
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
let ip = addr.ip();
let port = addr.port();
@@ -2747,8 +2755,12 @@ mod tests {
init_tracing(Level::ERROR);
let _ = rustfs_credentials::GLOBAL_RUSTFS_RPC_SECRET.set("test-rpc-secret".to_string());
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
let accept_task = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
drop(stream);
+14 -6
View File
@@ -558,16 +558,20 @@ mod tests {
use tokio::task::JoinHandle;
use tonic::transport::Endpoint as TonicEndpoint;
async fn spawn_hanging_listener() -> (String, JoinHandle<()>) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = format!("http://{}", listener.local_addr().unwrap());
async fn spawn_hanging_listener() -> Option<(String, JoinHandle<()>)> {
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = format!("http://{}", listener.local_addr().expect("listener local address should be available"));
let task = tokio::spawn(async move {
if let Ok((stream, _)) = listener.accept().await {
let _stream = stream;
tokio::time::sleep(Duration::from_secs(2)).await;
}
});
(addr, task)
Some((addr, task))
}
async fn cache_lazy_channel(addr: &str) {
@@ -588,7 +592,9 @@ mod tests {
#[tokio::test]
async fn test_remote_client_acquire_lock_respects_request_timeout_and_evicts_connection() {
ensure_test_rpc_secret();
let (addr, accept_task) = spawn_hanging_listener().await;
let Some((addr, accept_task)) = spawn_hanging_listener().await else {
return;
};
cache_lazy_channel(&addr).await;
assert!(GLOBAL_CONN_MAP.read().await.contains_key(&addr));
@@ -622,7 +628,9 @@ mod tests {
#[tokio::test]
async fn test_remote_client_acquire_locks_batch_respects_request_timeout_and_evicts_connection() {
ensure_test_rpc_secret();
let (addr, accept_task) = spawn_hanging_listener().await;
let Some((addr, accept_task)) = spawn_hanging_listener().await else {
return;
};
cache_lazy_channel(&addr).await;
assert!(GLOBAL_CONN_MAP.read().await.contains_key(&addr));
+540 -40
View File
@@ -18,12 +18,15 @@
use crate::batch_processor::{AsyncBatchProcessor, get_global_processors};
use crate::bitrot::{create_bitrot_reader, create_bitrot_writer};
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
use crate::bucket::metadata_sys;
use crate::bucket::object_lock::objectlock_sys::check_retention_for_modification;
use crate::bucket::replication::check_replicate_delete;
use crate::bucket::versioning::VersioningApi;
use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::client::{object_api_utils::get_raw_etag, transition_api::ReaderImpl};
use crate::disk::error_reduce::{OBJECT_OP_IGNORED_ERRS, reduce_read_quorum_errs, reduce_write_quorum_errs};
use crate::disk::error_reduce::{
BUCKET_OP_IGNORED_ERRS, OBJECT_OP_IGNORED_ERRS, count_errs, reduce_read_quorum_errs, reduce_write_quorum_errs,
};
use crate::disk::{
self, CHECK_PART_DISK_NOT_FOUND, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN,
conv_part_err_to_int, has_part_err,
@@ -34,6 +37,8 @@ use crate::error::{Error, Result, is_err_version_not_found};
use crate::error::{GenericError, ObjectApiError, is_err_object_not_found};
use crate::global::{GLOBAL_LocalNodeName, GLOBAL_TierConfigMgr};
use crate::object_api::ObjectOptions;
use crate::rpc::heal_bucket_local_on_disks;
use crate::store_utils::is_reserved_or_invalid_bucket;
use crate::{
bucket::lifecycle::bucket_lifecycle_ops::{
LifecycleOps, gen_transition_objname, get_transitioned_object_reader, put_restore_opts,
@@ -50,7 +55,7 @@ use crate::{
event_notification::{EventArgs, send_event},
global::{GLOBAL_LOCAL_DISK_MAP, GLOBAL_LOCAL_DISK_SET_DRIVES, get_global_deployment_id, is_dist_erasure},
object_api::{GetObjectReader, ObjectInfo, PutObjReader},
store_init::load_format_erasure,
store_init::{get_format_erasure_in_quorum, load_format_erasure, load_format_erasure_all, save_format_file},
};
use bytes::Bytes;
use bytesize::ByteSize;
@@ -75,7 +80,7 @@ use rustfs_lock::LockClient;
use rustfs_lock::fast_lock::types::LockResult;
use rustfs_lock::local_lock::LocalLock;
use rustfs_lock::{FastLockGuard, LockManager, NamespaceLock, NamespaceLockGuard, NamespaceLockWrapper, ObjectKey};
use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem};
use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos};
use rustfs_object_capacity::capacity_scope::{
CapacityScope, CapacityScopeDisk, record_capacity_scope, record_global_dirty_scope,
};
@@ -1772,23 +1777,205 @@ impl BucketOperations for SetDisks {
type Error = Error;
#[tracing::instrument(skip(self))]
async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> {
unimplemented!()
async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> {
let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1;
let force_create = opts.force_create;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks {
let bucket = bucket.to_string();
futures.push(async move {
match disk {
Some(disk) => match disk.make_volume(&bucket).await {
Ok(()) => Ok(()),
Err(err) if force_create && matches!(err, DiskError::VolumeExists) => Ok(()),
Err(err) => Err(err),
},
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let errs = results
.into_iter()
.map(|result| result.err())
.collect::<Vec<Option<DiskError>>>();
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
Ok(())
}
#[tracing::instrument(skip(self))]
async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result<BucketInfo> {
unimplemented!()
async fn get_bucket_info(&self, bucket: &str, _opts: &BucketOptions) -> Result<BucketInfo> {
let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks {
let bucket = bucket.to_string();
futures.push(async move {
match disk {
Some(disk) => disk.stat_volume(&bucket).await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut infos = Vec::with_capacity(results.len());
let mut errs = Vec::with_capacity(results.len());
for result in results {
match result {
Ok(info) => {
infos.push(Some(info));
errs.push(None);
}
Err(err) => {
infos.push(None);
errs.push(Some(err));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
let mut versioning = false;
let mut object_locking = false;
if let Ok(sys) = metadata_sys::get(bucket).await {
versioning = sys.versioning();
object_locking = sys.object_locking();
}
infos
.into_iter()
.flatten()
.next()
.map(|info| BucketInfo {
name: info.name,
created: info.created,
versioning,
object_locking,
..Default::default()
})
.ok_or(Error::VolumeNotFound)
}
#[tracing::instrument(skip(self))]
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
unimplemented!()
let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks {
futures.push(async move {
match disk {
Some(disk) => disk.list_volumes().await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut infos = Vec::with_capacity(results.len());
let mut errs = Vec::with_capacity(results.len());
for result in results {
match result {
Ok(volumes) => {
infos.push(Some(volumes));
errs.push(None);
}
Err(err) => {
infos.push(None);
errs.push(Some(err));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
let mut counts: HashMap<String, (usize, BucketInfo)> = HashMap::new();
for volumes in infos.into_iter().flatten() {
for volume in volumes {
if is_reserved_or_invalid_bucket(&volume.name, false) {
continue;
}
let entry = counts.entry(volume.name.clone()).or_insert((
0,
BucketInfo {
name: volume.name.clone(),
created: volume.created,
..Default::default()
},
));
entry.0 += 1;
}
}
let mut buckets = counts
.into_values()
.filter_map(|(count, bucket)| (count >= write_quorum).then_some(bucket))
.collect::<Vec<_>>();
buckets.sort_by(|left, right| left.name.cmp(&right.name));
Ok(buckets)
}
#[tracing::instrument(skip(self))]
async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> {
unimplemented!()
async fn delete_bucket(&self, bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> {
let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks.iter().cloned() {
let bucket = bucket.to_string();
futures.push(async move {
match disk {
Some(disk) => disk.delete_volume(&bucket).await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut errs = Vec::with_capacity(results.len());
let mut recreate = false;
for result in results {
match result {
Ok(()) => errs.push(None),
Err(err) => {
if matches!(err, DiskError::VolumeNotEmpty) {
recreate = true;
}
errs.push(Some(err));
}
}
}
if recreate {
for (index, err) in errs.iter().enumerate() {
if err.is_none()
&& let Some(Some(disk)) = disks.get(index)
{
let _ = disk.make_volume(bucket).await;
}
}
return Err(Error::VolumeNotEmpty);
}
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
Ok(())
}
}
@@ -3027,40 +3214,51 @@ impl rustfs_storage_api::ListOperations for SetDisks {
#[tracing::instrument(skip(self))]
async fn list_objects_v2(
self: Arc<Self>,
_bucket: &str,
_prefix: &str,
_continuation_token: Option<String>,
_delimiter: Option<String>,
_max_keys: i32,
_fetch_owner: bool,
_start_after: Option<String>,
_incl_deleted: bool,
bucket: &str,
prefix: &str,
continuation_token: Option<String>,
delimiter: Option<String>,
max_keys: i32,
fetch_owner: bool,
start_after: Option<String>,
incl_deleted: bool,
) -> Result<ListObjectsV2Info> {
unimplemented!()
self.inner_list_objects_v2(
bucket,
prefix,
continuation_token,
delimiter,
max_keys,
fetch_owner,
start_after,
incl_deleted,
)
.await
}
#[tracing::instrument(skip(self))]
async fn list_object_versions(
self: Arc<Self>,
_bucket: &str,
_prefix: &str,
_marker: Option<String>,
_version_marker: Option<String>,
_delimiter: Option<String>,
_max_keys: i32,
bucket: &str,
prefix: &str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
unimplemented!()
self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys)
.await
}
async fn walk(
self: Arc<Self>,
_rx: CancellationToken,
_bucket: &str,
_prefix: &str,
_result: Sender<ObjectInfoOrErr>,
_opts: WalkOptions,
rx: CancellationToken,
bucket: &str,
prefix: &str,
result: Sender<ObjectInfoOrErr>,
opts: WalkOptions,
) -> Result<()> {
unimplemented!()
self.walk_internal(rx, bucket, prefix, result, opts).await
}
}
@@ -3092,7 +3290,7 @@ impl rustfs_storage_api::MultipartOperations for SetDisks {
_src_opts: &ObjectOptions,
_dst_opts: &ObjectOptions,
) -> Result<()> {
unimplemented!()
Err(StorageError::NotImplemented)
}
#[tracing::instrument(level = "debug", skip(self, data, opts))]
@@ -4155,13 +4353,67 @@ impl rustfs_storage_api::HealOperations for SetDisks {
type HealOptions = HealOpts;
#[tracing::instrument(skip(self))]
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
unimplemented!()
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
let disks = self.disks.read().await.clone();
let (formats, errs) = load_format_erasure_all(&disks, true).await;
let ref_format = match get_format_erasure_in_quorum(&formats) {
Ok(format) => format,
Err(err) => {
let can_use_cached_layout = count_errs(&errs, &DiskError::UnformattedDisk) > 0
&& formats.iter().flatten().all(|format| self.format.check_other(format).is_ok())
&& errs
.iter()
.all(|err| err.is_none() || matches!(err, Some(DiskError::UnformattedDisk)));
if can_use_cached_layout {
self.format.clone()
} else {
return Ok((HealResultItem::default(), Some(err)));
}
}
};
let endpoints = crate::endpoints::Endpoints::from(self.set_endpoints.clone());
let before_drives = crate::layout::set_heal::formats_to_drives_info(&endpoints, &formats, &errs);
let mut result = HealResultItem {
heal_item_type: HealItemType::Metadata.to_string(),
detail: "disk-format".to_string(),
disk_count: self.set_drive_count,
set_count: 1,
before: Infos {
drives: before_drives.clone(),
},
after: Infos { drives: before_drives },
..Default::default()
};
if count_errs(&errs, &DiskError::UnformattedDisk) == 0 {
info!("set disk formats success, NoHealRequired, errs: {:?}", errs);
return Ok((result, Some(StorageError::NoHealRequired)));
}
if !dry_run {
for (disk_idx, err) in errs.iter().enumerate() {
if !matches!(err, Some(DiskError::UnformattedDisk)) {
continue;
}
let mut new_format = ref_format.clone();
new_format.erasure.this = ref_format.erasure.sets[self.set_index][disk_idx];
if save_format_file(&disks[disk_idx], &Some(new_format.clone())).await.is_ok() {
result.after.drives[disk_idx].uuid = new_format.erasure.this.to_string();
result.after.drives[disk_idx].state = DriveState::Ok.to_string();
}
}
}
Ok((result, None))
}
#[tracing::instrument(skip(self))]
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result<HealResultItem> {
unimplemented!()
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let mut result = heal_bucket_local_on_disks(bucket, opts, self.disk_inventory().await).await?;
result.set_count = 1;
Ok(result)
}
#[tracing::instrument(skip(self))]
@@ -4237,13 +4489,23 @@ impl rustfs_storage_api::HealOperations for SetDisks {
}
#[tracing::instrument(skip(self))]
async fn get_pool_and_set(&self, _id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)> {
unimplemented!()
async fn get_pool_and_set(&self, id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)> {
for (set_idx, set) in self.format.erasure.sets.iter().enumerate() {
for (disk_idx, disk_id) in set.iter().enumerate() {
if disk_id.to_string() == id {
return Ok((Some(self.pool_index), Some(set_idx), Some(disk_idx)));
}
}
}
Err(Error::DiskNotFound)
}
#[tracing::instrument(skip(self))]
async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> {
unimplemented!()
// Multipart orphan reconciliation is intentionally retained above the set layer
// until there is a concrete caller and a stable lower-level contract to implement.
Err(StorageError::NotImplemented)
}
}
@@ -5040,6 +5302,8 @@ mod tests {
use rustfs_filemeta::ReplicationState;
use rustfs_lock::client::local::LocalClient;
use rustfs_lock::{LockError, LockInfo, LockResponse, LockStats};
use rustfs_storage_api::HealOperations as _;
use rustfs_storage_api::ListOperations as _;
use rustfs_storage_api::TransitionedObject;
use rustfs_storage_api::{CompletePart, NamespaceLocking as _, ObjectOperations as _};
use serial_test::serial;
@@ -7256,4 +7520,240 @@ mod tests {
);
}
}
async fn make_local_bucket_test_set_disks() -> Arc<SetDisks> {
let format = FormatV3::new(1, 2);
let mut endpoints = Vec::new();
let mut disks = Vec::new();
for disk_idx in 0..2 {
let dir = tempfile::tempdir().expect("tempdir should be created");
let mut endpoint =
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
endpoint.set_pool_index(0);
endpoint.set_set_index(0);
endpoint.set_disk_index(disk_idx);
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("disk should be created");
let mut disk_format = format.clone();
disk_format.erasure.this = format.erasure.sets[0][disk_idx];
save_format_file(&Some(disk.clone()), &Some(disk_format))
.await
.expect("format should be saved");
std::mem::forget(dir);
endpoints.push(endpoint);
disks.push(Some(disk));
}
SetDisks::new(
"test-owner".to_string(),
Arc::new(RwLock::new(disks)),
2,
1,
0,
0,
endpoints,
format,
Vec::new(),
)
.await
}
async fn make_local_bucket_test_set_disks_with_missing_format() -> Arc<SetDisks> {
let format = FormatV3::new(1, 2);
let mut endpoints = Vec::new();
let mut disks = Vec::new();
for disk_idx in 0..2 {
let dir = tempfile::tempdir().expect("tempdir should be created");
let mut endpoint =
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
endpoint.set_pool_index(0);
endpoint.set_set_index(0);
endpoint.set_disk_index(disk_idx);
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("disk should be created");
if disk_idx == 0 {
let mut disk_format = format.clone();
disk_format.erasure.this = format.erasure.sets[0][disk_idx];
save_format_file(&Some(disk.clone()), &Some(disk_format))
.await
.expect("format should be saved");
}
std::mem::forget(dir);
endpoints.push(endpoint);
disks.push(Some(disk));
}
SetDisks::new(
"test-owner".to_string(),
Arc::new(RwLock::new(disks)),
2,
1,
0,
0,
endpoints,
format,
Vec::new(),
)
.await
}
#[tokio::test]
async fn bucket_operations_round_trip_without_panicking() {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "bucket-roundtrip";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let info = set_disks
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect("bucket info should be available");
assert_eq!(info.name, bucket);
let buckets = set_disks
.list_bucket(&BucketOptions::default())
.await
.expect("bucket listing should succeed");
assert!(buckets.iter().any(|entry| entry.name == bucket));
set_disks
.delete_bucket(bucket, &DeleteBucketOptions::default())
.await
.expect("bucket should be deleted");
}
#[tokio::test]
async fn set_level_listing_trait_methods_use_existing_listing_implementation() {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "bucket-listing";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(b"hello".to_vec());
set_disks
.put_object(bucket, "object", &mut reader, &ObjectOptions::default())
.await
.expect("object should be written");
let list_result = set_disks
.clone()
.list_objects_v2(bucket, "", None, None, 1000, false, None, false)
.await
.expect("set-level list_objects_v2 should succeed");
assert_eq!(list_result.objects.len(), 1);
assert_eq!(list_result.objects[0].name, "object");
let versions_result = set_disks
.clone()
.list_object_versions(bucket, "", None, None, None, 1000)
.await
.expect("set-level list_object_versions should succeed");
assert_eq!(versions_result.objects.len(), 1);
assert_eq!(versions_result.objects[0].name, "object");
let (tx, mut rx) = mpsc::channel(4);
set_disks
.clone()
.walk(CancellationToken::new(), bucket, "", tx, WalkOptions::default())
.await
.expect("set-level walk should succeed");
let mut walked_names = Vec::new();
while let Some(item) = rx.recv().await {
if let Some(object) = item.item {
walked_names.push(object.name);
}
}
assert!(walked_names.iter().any(|name| name == "object"));
}
#[tokio::test]
async fn set_level_heal_format_repairs_unformatted_disk() {
let set_disks = make_local_bucket_test_set_disks_with_missing_format().await;
let disk = {
let disks = set_disks.disks.read().await;
disks[1].clone().expect("second disk should exist")
};
let before = load_format_erasure(&disk, true)
.await
.expect_err("second disk should start unformatted");
assert_eq!(before, DiskError::UnformattedDisk);
let (heal_result, heal_err) = set_disks.heal_format(false).await.expect("heal_format should complete");
assert!(heal_err.is_none(), "heal_format should repair the local unformatted disk");
assert_eq!(heal_result.disk_count, 2);
assert_eq!(heal_result.set_count, 1);
assert_eq!(heal_result.after.drives[1].state, DriveState::Ok.to_string());
let repaired = load_format_erasure(&disk, true)
.await
.expect("second disk should contain a healed format");
assert_eq!(repaired.erasure.this, set_disks.format.erasure.sets[0][1]);
}
#[tokio::test]
async fn remaining_unsupported_trait_stubs_return_typed_errors() {
let set_disks = make_test_set_disks(Vec::new()).await;
let (heal_result, heal_err) = make_local_bucket_test_set_disks()
.await
.heal_format(false)
.await
.expect("heal_format should be callable on formatted disks");
assert!(matches!(heal_err, Some(StorageError::NoHealRequired)));
assert_eq!(heal_result.disk_count, 2);
let copy_part_err = set_disks
.copy_object_part(
"bucket",
"src",
"bucket",
"dst",
"upload-id",
1,
0,
1,
&ObjectInfo::default(),
&ObjectOptions::default(),
&ObjectOptions::default(),
)
.await
.expect_err("unsupported copy_object_part should return a typed error");
assert!(matches!(copy_part_err, StorageError::NotImplemented));
let abandoned_err = set_disks
.check_abandoned_parts("bucket", "object", &HealOpts::default())
.await
.expect_err("abandoned-parts check should stay in the upper reconciliation layer");
assert!(matches!(abandoned_err, StorageError::NotImplemented));
}
}
+278 -38
View File
@@ -440,22 +440,62 @@ impl BucketOperations for Sets {
type Error = Error;
#[tracing::instrument(skip(self))]
async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> {
unimplemented!()
async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> {
for set in &self.disk_set {
set.make_bucket(bucket, opts).await?;
}
Ok(())
}
#[tracing::instrument(skip(self))]
async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result<BucketInfo> {
unimplemented!()
async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo> {
let mut first_err = None;
for set in &self.disk_set {
match set.get_bucket_info(bucket, opts).await {
Ok(info) => return Ok(info),
Err(err) if first_err.is_none() => first_err = Some(err),
Err(_) => {}
}
}
Err(first_err.unwrap_or_else(|| StorageError::BucketNotFound(bucket.to_string())))
}
#[tracing::instrument(skip(self))]
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
unimplemented!()
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
let mut buckets = HashMap::new();
let mut first_err = None;
for set in &self.disk_set {
match set.list_bucket(opts).await {
Ok(set_buckets) => {
for bucket in set_buckets {
buckets.entry(bucket.name.clone()).or_insert(bucket);
}
}
Err(err) if first_err.is_none() => first_err = Some(err),
Err(_) => {}
}
}
if buckets.is_empty()
&& let Some(err) = first_err
{
return Err(err);
}
let mut buckets = buckets.into_values().collect::<Vec<_>>();
buckets.sort_by(|left, right| left.name.cmp(&right.name));
Ok(buckets)
}
#[tracing::instrument(skip(self))]
async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> {
unimplemented!()
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
for set in &self.disk_set {
set.delete_bucket(bucket, opts).await?;
}
Ok(())
}
}
@@ -543,8 +583,10 @@ impl rustfs_storage_api::ObjectOperations for Sets {
}
#[tracing::instrument(skip(self))]
async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, _force_del_marker: bool) -> Result<()> {
unimplemented!()
async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> {
self.get_disks_by_key(object)
.delete_object_version(bucket, object, fi, force_del_marker)
.await
}
#[tracing::instrument(skip(self))]
@@ -678,40 +720,51 @@ impl rustfs_storage_api::ListOperations for Sets {
#[tracing::instrument(skip(self))]
async fn list_objects_v2(
self: Arc<Self>,
_bucket: &str,
_prefix: &str,
_continuation_token: Option<String>,
_delimiter: Option<String>,
_max_keys: i32,
_fetch_owner: bool,
_start_after: Option<String>,
_incl_deleted: bool,
bucket: &str,
prefix: &str,
continuation_token: Option<String>,
delimiter: Option<String>,
max_keys: i32,
fetch_owner: bool,
start_after: Option<String>,
incl_deleted: bool,
) -> Result<ListObjectsV2Info> {
unimplemented!()
self.inner_list_objects_v2(
bucket,
prefix,
continuation_token,
delimiter,
max_keys,
fetch_owner,
start_after,
incl_deleted,
)
.await
}
#[tracing::instrument(skip(self))]
async fn list_object_versions(
self: Arc<Self>,
_bucket: &str,
_prefix: &str,
_marker: Option<String>,
_version_marker: Option<String>,
_delimiter: Option<String>,
_max_keys: i32,
bucket: &str,
prefix: &str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
unimplemented!()
self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys)
.await
}
async fn walk(
self: Arc<Self>,
_rx: CancellationToken,
_bucket: &str,
_prefix: &str,
_result: tokio::sync::mpsc::Sender<ObjectInfoOrErr>,
_opts: WalkOptions,
rx: CancellationToken,
bucket: &str,
prefix: &str,
result: tokio::sync::mpsc::Sender<ObjectInfoOrErr>,
opts: WalkOptions,
) -> Result<()> {
unimplemented!()
self.walk_internal(rx, bucket, prefix, result, opts).await
}
}
@@ -762,7 +815,7 @@ impl rustfs_storage_api::MultipartOperations for Sets {
_src_opts: &ObjectOptions,
_dst_opts: &ObjectOptions,
) -> Result<()> {
unimplemented!()
Err(StorageError::NotImplemented)
}
#[tracing::instrument(skip(self))]
@@ -920,8 +973,22 @@ impl rustfs_storage_api::HealOperations for Sets {
Ok((res, None))
}
#[tracing::instrument(skip(self))]
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result<HealResultItem> {
unimplemented!()
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let mut result = HealResultItem {
heal_item_type: HealItemType::Bucket.to_string(),
bucket: bucket.to_string(),
set_count: self.set_count,
..Default::default()
};
for set in &self.disk_set {
let mut set_result = set.heal_bucket(bucket, opts).await?;
result.disk_count += set_result.disk_count;
result.before.drives.append(&mut set_result.before.drives);
result.after.drives.append(&mut set_result.after.drives);
}
Ok(result)
}
#[tracing::instrument(skip(self))]
async fn heal_object(
@@ -936,12 +1003,22 @@ impl rustfs_storage_api::HealOperations for Sets {
.await
}
#[tracing::instrument(skip(self))]
async fn get_pool_and_set(&self, _id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)> {
unimplemented!()
async fn get_pool_and_set(&self, id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)> {
for (set_idx, set) in self.format.erasure.sets.iter().enumerate() {
for (disk_idx, disk_id) in set.iter().enumerate() {
if disk_id.to_string() == id {
return Ok((Some(self.pool_idx), Some(set_idx), Some(disk_idx)));
}
}
}
Err(Error::DiskNotFound)
}
#[tracing::instrument(skip(self))]
async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> {
unimplemented!()
// Multipart orphan reconciliation is intentionally retained above the pool/set layers
// until there is a concrete caller and a stable lower-level contract to implement.
Err(StorageError::NotImplemented)
}
}
@@ -1021,6 +1098,9 @@ async fn init_storage_disks_with_errors(
#[cfg(test)]
mod tests {
use super::*;
use crate::layout::endpoint::Endpoint;
use rustfs_storage_api::HealOperations as _;
use rustfs_storage_api::ListOperations as _;
#[test]
fn test_apply_delete_objects_results_preserves_original_order_for_out_of_order_batches() {
@@ -1089,4 +1169,164 @@ mod tests {
Some(Error::other("third failed").to_string())
);
}
#[tokio::test]
async fn sets_get_pool_and_set_returns_matching_coordinates() {
let format = FormatV3::new(2, 2);
let target = format.erasure.sets[1][0].to_string();
let endpoints = vec![
Endpoint::try_from("http://127.0.0.1:9000/data0").expect("first endpoint should parse"),
Endpoint::try_from("http://127.0.0.1:9001/data1").expect("second endpoint should parse"),
Endpoint::try_from("http://127.0.0.1:9002/data2").expect("third endpoint should parse"),
Endpoint::try_from("http://127.0.0.1:9003/data3").expect("fourth endpoint should parse"),
];
let sets = Sets {
id: format.id,
disk_set: Vec::new(),
pool_idx: 3,
endpoints: PoolEndpoints {
legacy: false,
set_count: 2,
drives_per_set: 2,
endpoints: Endpoints::from(endpoints),
cmd_line: String::new(),
platform: String::new(),
},
format,
parity_count: 1,
set_count: 2,
set_drive_count: 2,
default_parity_count: 1,
distribution_algo: DistributionAlgoVersion::V1,
exit_signal: None,
};
let result = sets
.get_pool_and_set(&target)
.await
.expect("disk id should resolve within the pool");
assert_eq!(result, (Some(3), Some(1), Some(0)));
}
#[tokio::test]
async fn sets_list_objects_v2_lists_objects_within_the_pool() {
let format = FormatV3::new(1, 2);
let mut endpoints = Vec::new();
let mut disks = Vec::new();
for disk_idx in 0..2 {
let dir = tempfile::tempdir().expect("tempdir should be created");
let mut endpoint =
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
endpoint.set_pool_index(0);
endpoint.set_set_index(0);
endpoint.set_disk_index(disk_idx);
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("disk should be created");
let mut disk_format = format.clone();
disk_format.erasure.this = format.erasure.sets[0][disk_idx];
save_format_file(&Some(disk.clone()), &Some(disk_format))
.await
.expect("format should be saved");
std::mem::forget(dir);
endpoints.push(endpoint);
disks.push(Some(disk));
}
let set_disks = SetDisks::new(
"test-owner".to_string(),
Arc::new(RwLock::new(disks)),
2,
1,
0,
0,
endpoints.clone(),
format.clone(),
Vec::new(),
)
.await;
let sets = Arc::new(Sets {
id: format.id,
disk_set: vec![set_disks],
pool_idx: 0,
endpoints: PoolEndpoints {
legacy: false,
set_count: 1,
drives_per_set: 2,
endpoints: Endpoints::from(endpoints),
cmd_line: String::new(),
platform: String::new(),
},
format,
parity_count: 1,
set_count: 1,
set_drive_count: 2,
default_parity_count: 1,
distribution_algo: DistributionAlgoVersion::V1,
exit_signal: None,
});
sets.make_bucket("bucket", &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(b"hello".to_vec());
sets.put_object("bucket", "object", &mut reader, &ObjectOptions::default())
.await
.expect("object should be written");
let result = sets
.clone()
.list_objects_v2("bucket", "", None, None, 1000, false, None, false)
.await
.expect("pool-level listing should succeed");
assert_eq!(result.objects.len(), 1);
assert_eq!(result.objects[0].name, "object");
}
#[tokio::test]
async fn sets_check_abandoned_parts_returns_typed_not_implemented_error() {
let format = FormatV3::new(1, 1);
let sets = Sets {
id: format.id,
disk_set: Vec::new(),
pool_idx: 0,
endpoints: PoolEndpoints {
legacy: false,
set_count: 1,
drives_per_set: 1,
endpoints: Endpoints::from(Vec::new()),
cmd_line: String::new(),
platform: String::new(),
},
format,
parity_count: 0,
set_count: 1,
set_drive_count: 1,
default_parity_count: 0,
distribution_algo: DistributionAlgoVersion::V1,
exit_signal: None,
};
let err = sets
.check_abandoned_parts("bucket", "object", &HealOpts::default())
.await
.expect_err("abandoned-parts ownership should stay above the pool/set storage layers");
assert!(matches!(err, StorageError::NotImplemented));
}
}
+7 -18
View File
@@ -156,23 +156,12 @@ impl ECStore {
#[instrument(skip(self))]
pub(super) async fn handle_check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> {
let object = encode_dir_object(object);
if self.single_pool() {
return self.pools[0].check_abandoned_parts(bucket, &object, opts).await;
}
let mut errs = Vec::new();
for pool in self.pools.iter() {
//TODO: IsSuspended
if let Err(err) = pool.check_abandoned_parts(bucket, &object, opts).await {
errs.push(err);
}
}
if !errs.is_empty() {
return Err(errs[0].clone());
}
Ok(())
let _ = (bucket, object, opts);
// Stale multipart reconciliation is already owned by the lifecycle-driven
// background cleanup path in `bucket_lifecycle_ops.rs`. There is currently
// no stable object-heal contract that should fan this request out through
// pool/set storage layers, so keep the placeholder explicit at the ECStore
// boundary instead of dispatching into lower layers.
Err(StorageError::NotImplemented)
}
}
+2 -3
View File
@@ -180,9 +180,8 @@ impl ECStore {
) -> Result<()> {
check_new_multipart_args(src_bucket, src_object)?;
// TODO: PutObjectReader
// self.put_object_part(dst_bucket, dst_object, upload_id, part_id, data, opts)
// The full UploadPartCopy path still requires the higher S3/request layer to
// derive encryption, compression, and multipart checksum write semantics.
Err(StorageError::NotImplemented)
}
File diff suppressed because it is too large Load Diff
+14 -6
View File
@@ -1634,7 +1634,7 @@ mod tests {
fn start_mock_oidc_discovery_server<F>(
build_discovery_issuer: F,
max_requests: usize,
) -> (String, std::thread::JoinHandle<()>)
) -> Option<(String, std::thread::JoinHandle<()>)>
where
F: Fn(&str) -> String + Send + 'static,
{
@@ -1650,8 +1650,12 @@ mod tests {
const IDLE_SHUTDOWN: Duration = Duration::from_secs(1);
const ABSOLUTE_CAP: Duration = Duration::from_secs(5);
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let base = format!("http://{}", listener.local_addr().unwrap());
let listener = match TcpListener::bind("127.0.0.1:0") {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None,
Err(err) => panic!("test listener should bind: {err}"),
};
let base = format!("http://{}", listener.local_addr().expect("listener local address should be available"));
let discovery_issuer = build_discovery_issuer(&base);
let discovery_body = serde_json::json!({
"issuer": discovery_issuer,
@@ -1754,7 +1758,7 @@ mod tests {
.recv_timeout(Duration::from_millis(100))
.expect("mock OIDC discovery server should become ready");
(base, handle)
Some((base, handle))
}
fn discovery_error_contains_all_variants(err: &str, base: &str) -> bool {
@@ -1776,7 +1780,9 @@ mod tests {
async fn test_validate_oidc_provider_config_retries_with_issuer_candidates() {
// Discovery document must advertise the canonical issuer path. The first candidate has no
// trailing slash; openidconnect rejects issuer mismatch, then the second variant succeeds.
let (base, handle) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/rustfs/"), 8);
let Some((base, handle)) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/rustfs/"), 8) else {
return;
};
let config_url = format!("{base}/application/o/rustfs");
let config = build_mocked_oidc_provider_config("default", &config_url);
@@ -1789,7 +1795,9 @@ mod tests {
#[tokio::test]
async fn test_validate_oidc_provider_config_returns_detailed_errors() {
let (base, handle) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/other"), 8);
let Some((base, handle)) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/other"), 8) else {
return;
};
let config_url = format!("{base}/application/o/rustfs");
let config = build_mocked_oidc_provider_config("default", &config_url);
+29 -11
View File
@@ -1016,9 +1016,13 @@ mod tests {
StatusCode::OK
}
async fn start_test_server(state: TestState) -> (String, tokio::task::JoinHandle<()>) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
async fn start_test_server(state: TestState) -> Option<(String, tokio::task::JoinHandle<()>)> {
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
let app = Router::new()
.route("/stream", get(get_stream).head(reject_head).put(accept_put))
.route("/stall", get(get_stalling_stream))
@@ -1028,7 +1032,7 @@ mod tests {
axum::serve(listener, app).await.unwrap();
});
(format!("http://{addr}/stream"), handle)
Some((format!("http://{addr}/stream"), handle))
}
#[test]
@@ -1055,7 +1059,9 @@ mod tests {
#[tokio::test]
async fn http_reader_does_not_send_preflight_head() {
let state = TestState::default();
let (url, handle) = start_test_server(state.clone()).await;
let Some((url, handle)) = start_test_server(state.clone()).await else {
return;
};
let mut reader = HttpReader::new(url, Method::GET, HeaderMap::new(), None).await.unwrap();
let mut buf = Vec::new();
@@ -1071,7 +1077,9 @@ mod tests {
#[tokio::test]
async fn http_reader_stall_timeout_triggers_after_progress_stops() {
let state = TestState::default();
let (base_url, handle) = start_test_server(state.clone()).await;
let Some((base_url, handle)) = start_test_server(state.clone()).await else {
return;
};
let url = base_url.replace("/stream", "/stall");
let mut reader =
@@ -1099,7 +1107,9 @@ mod tests {
#[tokio::test]
async fn http_writer_does_not_send_empty_preflight_put() {
let state = TestState::default();
let (url, handle) = start_test_server(state.clone()).await;
let Some((url, handle)) = start_test_server(state.clone()).await else {
return;
};
let mut writer = HttpWriter::new(url, Method::PUT, HeaderMap::new()).await.unwrap();
writer.write_all(b"payload").await.unwrap();
@@ -1114,7 +1124,9 @@ mod tests {
#[tokio::test]
async fn http_writer_handles_many_small_writes() {
let state = TestState::default();
let (url, handle) = start_test_server(state.clone()).await;
let Some((url, handle)) = start_test_server(state.clone()).await else {
return;
};
let mut writer = HttpWriter::new(url, Method::PUT, HeaderMap::new()).await.unwrap();
let chunk = b"0123456789abcdef";
@@ -1134,7 +1146,9 @@ mod tests {
#[tokio::test]
async fn http_writer_supports_vectored_writes() {
let state = TestState::default();
let (url, handle) = start_test_server(state.clone()).await;
let Some((url, handle)) = start_test_server(state.clone()).await else {
return;
};
let mut writer = HttpWriter::new(url, Method::PUT, HeaderMap::new()).await.unwrap();
let bufs = [IoSlice::new(b"hello "), IoSlice::new(b"world")];
@@ -1150,8 +1164,12 @@ mod tests {
#[tokio::test]
async fn http_reader_request_error_includes_method_and_url() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
drop(listener);
let url = format!("http://{addr}/stream");
+25 -5
View File
@@ -1005,7 +1005,11 @@ mod tests {
#[tokio::test]
async fn is_active_succeeds_when_ping_returns_pong() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis");
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("bind fake redis: {err}"),
};
let addr = listener.local_addr().expect("listener addr");
tokio::spawn(run_fake_redis_server(listener, false));
@@ -1021,7 +1025,11 @@ mod tests {
#[tokio::test]
async fn is_active_returns_error_when_ping_fails() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis");
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("bind fake redis: {err}"),
};
let addr = listener.local_addr().expect("listener addr");
tokio::spawn(async move {
loop {
@@ -1165,7 +1173,11 @@ mod tests {
#[tokio::test]
async fn send_body_keeps_connected_true_when_retryable_error_eventually_recovers() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis");
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("bind fake redis: {err}"),
};
let addr = listener.local_addr().expect("listener addr");
tokio::spawn(run_fake_redis_server(listener, true));
@@ -1197,7 +1209,11 @@ mod tests {
#[tokio::test]
async fn send_body_sets_connected_false_after_retry_exhaustion() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis");
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("bind fake redis: {err}"),
};
let addr = listener.local_addr().expect("listener addr");
tokio::spawn(async move {
loop {
@@ -1238,7 +1254,11 @@ mod tests {
#[tokio::test]
async fn send_raw_from_store_failure_does_not_count_as_success() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis");
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("bind fake redis: {err}"),
};
let addr = listener.local_addr().expect("listener addr");
tokio::spawn(async move {
loop {
+50 -8
View File
@@ -265,12 +265,42 @@ pub fn get_available_port() -> u16 {
}
fn try_get_available_port() -> std::io::Result<u16> {
let listener =
TcpListener::bind("0.0.0.0:0").map_err(|err| Error::other(format!("Failed to bind for ephemeral port: {err}")))?;
listener
.local_addr()
.map(|addr| addr.port())
.map_err(|err| Error::other(format!("Failed to read ephemeral port: {err}")))
let mut last_err = None;
for _ in 0..8 {
for candidate in ["127.0.0.1:0", "0.0.0.0:0"] {
match TcpListener::bind(candidate) {
Ok(listener) => {
return listener
.local_addr()
.map(|addr| addr.port())
.map_err(|err| Error::new(err.kind(), format!("Failed to read ephemeral port: {err}")));
}
Err(err)
if matches!(
err.kind(),
std::io::ErrorKind::AddrInUse
| std::io::ErrorKind::AddrNotAvailable
| std::io::ErrorKind::Interrupted
| std::io::ErrorKind::WouldBlock
| std::io::ErrorKind::PermissionDenied
) =>
{
last_err = Some(err);
}
Err(err) => {
return Err(Error::new(err.kind(), format!("Failed to bind for ephemeral port on {candidate}: {err}")));
}
}
}
std::thread::sleep(Duration::from_millis(5));
}
match last_err {
Some(err) => Err(Error::new(err.kind(), format!("Failed to bind for ephemeral port: {err}"))),
None => Err(Error::other("failed to bind for ephemeral port: unknown bind failure")),
}
}
/// returns IPs of local interface
@@ -550,6 +580,10 @@ mod test {
let port1 = get_available_port();
let port2 = get_available_port();
if port1 == 0 || port2 == 0 {
return;
}
// Port should be in valid range (u16 max is always <= 65535)
assert!(port1 > 0);
assert!(port2 > 0);
@@ -656,7 +690,11 @@ mod test {
assert_eq!(result.port(), 8080);
// Test port-only format with port 0 (should get available port)
let result = parse_and_resolve_address(":0").unwrap();
let result = match parse_and_resolve_address(":0") {
Ok(result) => result,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("expected dynamic port resolution for :0: {err}"),
};
assert_eq!(result.ip(), IpAddr::V6(Ipv6Addr::UNSPECIFIED));
assert!(result.port() > 0);
@@ -665,7 +703,11 @@ mod test {
assert_eq!(result.port(), 9000);
// Test localhost with port 0 (should get available port)
let result = parse_and_resolve_address("localhost:0").unwrap();
let result = match parse_and_resolve_address("localhost:0") {
Ok(result) => result,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("expected dynamic port resolution for localhost:0: {err}"),
};
assert!(result.port() > 0);
// Test 0.0.0.0 with port
+31 -6
View File
@@ -22,10 +22,12 @@ use rustfs_common::{GlobalReadiness, set_global_addr};
use rustfs_credentials::init_global_action_credentials;
use rustfs_utils::net::parse_and_resolve_address;
use std::{
io::{Error, Result},
io::{Error, ErrorKind, Result},
net::SocketAddr,
path::Path,
sync::Arc,
thread,
time::Duration,
};
use tempfile::TempDir;
use tracing::{debug, error, info, warn};
@@ -165,10 +167,29 @@ pub(crate) async fn prepare_embedded_startup_config(
}
pub(crate) fn find_embedded_available_port() -> Result<u16> {
let listener = std::net::TcpListener::bind("127.0.0.1:0")?;
let port = listener.local_addr()?.port();
drop(listener);
Ok(port)
let mut last_err = None;
for _ in 0..8 {
match std::net::TcpListener::bind("127.0.0.1:0") {
Ok(listener) => {
let port = listener.local_addr()?.port();
drop(listener);
return Ok(port);
}
Err(err)
if matches!(
err.kind(),
ErrorKind::AddrInUse | ErrorKind::AddrNotAvailable | ErrorKind::Interrupted | ErrorKind::WouldBlock
) =>
{
last_err = Some(err);
thread::sleep(Duration::from_millis(5));
}
Err(err) => return Err(err),
}
}
Err(last_err.unwrap_or_else(|| Error::other("failed to reserve an embedded TCP port")))
}
pub(crate) async fn init_embedded_startup_listen_context(config: &Config) -> Result<EmbeddedStartupListenContext> {
@@ -374,7 +395,11 @@ mod tests {
#[test]
fn find_embedded_available_port_returns_tcp_port() {
let port = find_embedded_available_port().expect("available port should be found");
let port = match find_embedded_available_port() {
Ok(port) => port,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("available port should be found: {err}"),
};
assert_ne!(port, 0);
}
+18 -6
View File
@@ -2813,9 +2813,13 @@ mod tests {
);
}
async fn connect_test_node_service_client() -> NodeServiceClient<tonic::transport::Channel> {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
async fn connect_test_node_service_client() -> Option<NodeServiceClient<tonic::transport::Channel>> {
let listener = match TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
let service = create_test_node_service();
tokio::spawn(async move {
@@ -2826,12 +2830,18 @@ mod tests {
.unwrap();
});
NodeServiceClient::connect(format!("http://{addr}")).await.unwrap()
Some(
NodeServiceClient::connect(format!("http://{addr}"))
.await
.expect("node service test client should connect"),
)
}
#[tokio::test]
async fn test_write_stream_unimplemented() {
let mut client = connect_test_node_service_client().await;
let Some(mut client) = connect_test_node_service_client().await else {
return;
};
let request = tokio_stream::iter([WriteRequest::default()]);
let response = client.write_stream(request).await;
@@ -2843,7 +2853,9 @@ mod tests {
#[tokio::test]
async fn test_read_at_unimplemented() {
let mut client = connect_test_node_service_client().await;
let Some(mut client) = connect_test_node_service_client().await else {
return;
};
let request = tokio_stream::iter([ReadAtRequest::default()]);
let response = client.read_at(request).await;
+5 -1
View File
@@ -45,7 +45,11 @@ async fn test_embedded_server_recovers_after_deferred_iam_bootstrap() {
(ENV_TEST_IAM_RETRY_INTERVAL_MS, Some("500")),
],
async {
let port = find_available_port().expect("find free port");
let port = match find_available_port() {
Ok(port) => port,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("find free port: {err}"),
};
let server = RustFSServerBuilder::new()
.address(format!("127.0.0.1:{port}"))
.access_key("testaccesskey")
+5 -1
View File
@@ -38,7 +38,11 @@ fn s3_client(endpoint: &str, access_key: &str, secret_key: &str) -> Client {
#[tokio::test]
async fn test_embedded_server_basic_s3_operations() {
// 1. Pick a free port and start the embedded server.
let port = find_available_port().expect("find free port");
let port = match find_available_port() {
Ok(port) => port,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("find free port: {err}"),
};
let server = RustFSServerBuilder::new()
.address(format!("127.0.0.1:{port}"))
.access_key("testaccesskey")
+1
View File
@@ -50,6 +50,7 @@ accepted|rustfs/src/server/readiness.rs|infra->app|crate::app::context::resolve_
accepted|rustfs/src/startup_iam.rs|infra->app|crate::app::context::AppContext|startup wires IAM bootstrap through AppContext
accepted|rustfs/src/storage/ecfs_extend.rs|infra->app|crate::app::context::resolve_buffer_config|storage buffer sizing reads runtime config through AppContext resolver
accepted|rustfs/src/storage/ecfs_extend.rs|infra->interface|crate::storage::ecfs::ListObjectUnorderedQuery|storage extension uses current ECFS interface query type
accepted|rustfs/src/storage/ecfs_extend.rs|infra->app|crate::app::context::resolve_buffer_config|buffer config resolution uses global AppContext
accepted|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::ecfs::FS|storage tests exercise current ECFS interface path
accepted|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::ecfs::validate_object_lock_configuration_input|storage tests exercise current ECFS object-lock validator
accepted|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::s3_api::common::rustfs_initiator|storage tests use current S3 API initiator helper