fix scstore::new bugs, stop ns_lock

This commit is contained in:
weisd
2025-04-19 01:52:27 +08:00
parent cf70e12e14
commit 8ea6f7e627
2 changed files with 63 additions and 72 deletions
+53 -58
View File
@@ -59,12 +59,7 @@ use common::error::{Error, Result};
use futures::future::join_all; use futures::future::join_all;
use glob::Pattern; use glob::Pattern;
use http::HeaderMap; use http::HeaderMap;
use lock::{ use lock::{namespace_lock::NsLockMap, LockApi};
// drwmutex::Options,
drwmutex::Options,
namespace_lock::{new_nslock, NsLockMap},
LockApi,
};
use madmin::heal_commands::{HealDriveInfo, HealResultItem}; use madmin::heal_commands::{HealDriveInfo, HealResultItem};
use md5::{Digest as Md5Digest, Md5}; use md5::{Digest as Md5Digest, Md5};
use rand::{ use rand::{
@@ -3672,33 +3667,33 @@ impl ObjectIO for SetDisks {
async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> { async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
let disks = self.disks.read().await; let disks = self.disks.read().await;
let mut _ns = None; // let mut _ns = None;
if !opts.no_lock { // if !opts.no_lock {
let paths = vec![object.to_string()]; // let paths = vec![object.to_string()];
let ns_lock = new_nslock( // let ns_lock = new_nslock(
Arc::clone(&self.ns_mutex), // Arc::clone(&self.ns_mutex),
self.locker_owner.clone(), // self.locker_owner.clone(),
bucket.to_string(), // bucket.to_string(),
paths, // paths,
self.lockers.clone(), // self.lockers.clone(),
) // )
.await; // .await;
if !ns_lock // if !ns_lock
.0 // .0
.write() // .write()
.await // .await
.get_lock(&Options { // .get_lock(&Options {
timeout: Duration::from_secs(5), // timeout: Duration::from_secs(5),
retry_interval: Duration::from_secs(1), // retry_interval: Duration::from_secs(1),
}) // })
.await // .await
.map_err(|err| Error::from_string(err.to_string()))? // .map_err(|err| Error::from_string(err.to_string()))?
{ // {
return Err(Error::from_string("can not get lock. please retry".to_string())); // return Err(Error::from_string("can not get lock. please retry".to_string()));
} // }
_ns = Some(ns_lock); // _ns = Some(ns_lock);
} // }
let mut user_defined = opts.user_defined.clone().unwrap_or_default(); let mut user_defined = opts.user_defined.clone().unwrap_or_default();
@@ -4165,33 +4160,33 @@ impl StorageAPI for SetDisks {
unimplemented!() unimplemented!()
} }
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> { async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> {
let mut _ns = None; // let mut _ns = None;
if !opts.no_lock { // if !opts.no_lock {
let paths = vec![object.to_string()]; // let paths = vec![object.to_string()];
let ns_lock = new_nslock( // let ns_lock = new_nslock(
Arc::clone(&self.ns_mutex), // Arc::clone(&self.ns_mutex),
self.locker_owner.clone(), // self.locker_owner.clone(),
bucket.to_string(), // bucket.to_string(),
paths, // paths,
self.lockers.clone(), // self.lockers.clone(),
) // )
.await; // .await;
if !ns_lock // if !ns_lock
.0 // .0
.write() // .write()
.await // .await
.get_lock(&Options { // .get_lock(&Options {
timeout: Duration::from_secs(5), // timeout: Duration::from_secs(5),
retry_interval: Duration::from_secs(1), // retry_interval: Duration::from_secs(1),
}) // })
.await // .await
.map_err(|err| Error::from_string(err.to_string()))? // .map_err(|err| Error::from_string(err.to_string()))?
{ // {
return Err(Error::from_string("can not get lock. please retry".to_string())); // return Err(Error::from_string("can not get lock. please retry".to_string()));
} // }
_ns = Some(ns_lock); // _ns = Some(ns_lock);
} // }
let (fi, _, _) = self.get_object_fileinfo(bucket, object, opts, false).await?; let (fi, _, _) = self.get_object_fileinfo(bucket, object, opts, false).await?;
// warn!("get object_info fi {:?}", &fi); // warn!("get object_info fi {:?}", &fi);
+10 -14
View File
@@ -521,28 +521,24 @@ impl ECStore {
} }
} }
let mut server_pools = Vec::new(); let mut server_pools = vec![PoolAvailableSpace::default(); self.pools.len()];
for (i, zinfo) in infos.iter().enumerate() { for (i, zinfo) in infos.iter().enumerate() {
if zinfo.is_empty() { if zinfo.is_empty() {
server_pools.push(PoolAvailableSpace { server_pools[i] = PoolAvailableSpace {
index: i, index: i,
..Default::default() ..Default::default()
}); };
continue; continue;
} }
if !is_meta_bucketname(bucket) { if !is_meta_bucketname(bucket) && !has_space_for(zinfo, size).await.unwrap_or_default() {
let avail = has_space_for(zinfo, size).await.unwrap_or_default(); server_pools[i] = PoolAvailableSpace {
index: i,
..Default::default()
};
if !avail { continue;
server_pools.push(PoolAvailableSpace {
index: i,
..Default::default()
});
continue;
}
} }
let mut available = 0; let mut available = 0;
@@ -2530,7 +2526,7 @@ async fn get_disk_infos(disks: &[Option<DiskStore>]) -> Vec<Option<DiskInfo>> {
res res
} }
#[derive(Debug, Default)] #[derive(Debug, Default, Clone)]
pub struct PoolAvailableSpace { pub struct PoolAvailableSpace {
pub index: usize, pub index: usize,
pub available: u64, // in bytes pub available: u64, // in bytes