Merge pull request #163 from rustfs/fix/disks_join

Fix/disks join
This commit is contained in:
weisd
2024-12-12 09:57:43 +08:00
committed by GitHub
9 changed files with 195 additions and 165 deletions
+4
View File
@@ -265,6 +265,10 @@ pub fn os_err_to_file_err(e: io::Error) -> Error {
} }
} }
pub fn is_unformatted_disk(err: &Error) -> bool {
matches!(err.downcast_ref::<DiskError>(), Some(DiskError::UnformattedDisk))
}
pub fn is_err_file_not_found(err: &Error) -> bool { pub fn is_err_file_not_found(err: &Error) -> bool {
matches!(err.downcast_ref::<DiskError>(), Some(DiskError::FileNotFound)) matches!(err.downcast_ref::<DiskError>(), Some(DiskError::FileNotFound))
} }
+2
View File
@@ -247,6 +247,7 @@ impl LocalDisk {
true true
} }
#[tracing::instrument(level = "debug", skip(self))]
async fn check_format_json(&self) -> Result<Metadata> { async fn check_format_json(&self) -> Result<Metadata> {
let md = fs::metadata(&self.format_path).await.map_err(|e| match e.kind() { let md = fs::metadata(&self.format_path).await.map_err(|e| match e.kind() {
ErrorKind::NotFound => DiskError::DiskNotFound, ErrorKind::NotFound => DiskError::DiskNotFound,
@@ -801,6 +802,7 @@ impl DiskAPI for LocalDisk {
} }
} }
#[tracing::instrument(level = "debug", skip(self))]
async fn get_disk_id(&self) -> Result<Option<Uuid>> { async fn get_disk_id(&self) -> Result<Option<Uuid>> {
let mut format_info = self.format_info.write().await; let mut format_info = self.format_info.write().await;
+14 -7
View File
@@ -533,14 +533,21 @@ impl ShardReader {
let mut ress = Vec::with_capacity(reader_length); let mut ress = Vec::with_capacity(reader_length);
for disk in self.readers.iter_mut() { for disk in self.readers.iter_mut() {
if disk.is_none() { // if disk.is_none() {
ress.push(None); // ress.push(None);
errors.push(Some(Error::new(DiskError::DiskNotFound))); // errors.push(Some(Error::new(DiskError::DiskNotFound)));
continue; // continue;
} // }
let disk: &mut BitrotReader = disk.as_mut().unwrap(); // let disk: &mut BitrotReader = disk.as_mut().unwrap();
futures.push(disk.read_at(self.offset, read_length)); let offset = self.offset;
futures.push(async move {
if let Some(disk) = disk {
disk.read_at(offset, read_length).await
} else {
Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
+142 -144
View File
@@ -194,25 +194,24 @@ impl SetDisks {
let mut errs = Vec::with_capacity(disks.len()); let mut errs = Vec::with_capacity(disks.len());
for (i, disk) in disks.iter().enumerate() { for (i, disk) in disks.iter().enumerate() {
if disk.is_none() {
// ress.push(None);
errs.push(Some(Error::new(DiskError::DiskNotFound)));
continue;
}
let disk = disk.as_ref().unwrap();
let mut file_info = file_infos[i].clone(); let mut file_info = file_infos[i].clone();
if file_info.erasure.index == 0 { futures.push(async move {
file_info.erasure.index = i + 1; if file_info.erasure.index == 0 {
} file_info.erasure.index = i + 1;
}
if !file_info.is_valid() { if !file_info.is_valid() {
// ress.push(None); return Err(Error::new(DiskError::FileCorrupt));
errs.push(Some(Error::new(DiskError::FileCorrupt))); }
continue;
}
futures.push(disk.rename_data(src_bucket, src_object, file_info, dst_bucket, dst_object)) if let Some(disk) = disk {
disk.rename_data(src_bucket, src_object, file_info, dst_bucket, dst_object)
.await
} else {
Err(Error::new(DiskError::DiskNotFound))
}
})
} }
let mut disk_versions = vec![None; disks.len()]; let mut disk_versions = vec![None; disks.len()];
@@ -326,20 +325,22 @@ impl SetDisks {
let mut errs = Vec::with_capacity(disks.len()); let mut errs = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { let file_path = file_path.clone();
errs.push(Some(Error::new(DiskError::DiskNotFound))); futures.push(async move {
continue; if let Some(disk) = disk {
} disk.delete(
bucket,
let disk = disk.as_ref().unwrap(); &file_path,
futures.push(disk.delete( DeleteOptions {
bucket, recursive: true,
&file_path, ..Default::default()
DeleteOptions { },
recursive: true, )
..Default::default() .await
}, } else {
)); Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -367,12 +368,13 @@ impl SetDisks {
let mut errs = Vec::with_capacity(disks.len()); let mut errs = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { futures.push(async move {
errs.push(Some(Error::new(DiskError::DiskNotFound))); if let Some(disk) = disk {
continue; disk.delete_paths(RUSTFS_META_MULTIPART_BUCKET, paths).await
} } else {
let disk = disk.as_ref().unwrap(); Err(Error::new(DiskError::DiskNotFound))
futures.push(disk.delete_paths(RUSTFS_META_MULTIPART_BUCKET, paths)) }
})
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -401,12 +403,14 @@ impl SetDisks {
let mut errs = Vec::with_capacity(disks.len()); let mut errs = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { let meta = meta.clone();
errs.push(Some(Error::new(DiskError::DiskNotFound))); futures.push(async move {
continue; if let Some(disk) = disk {
} disk.rename_part(src_bucket, src_object, dst_bucket, dst_object, meta).await
let disk = disk.as_ref().unwrap(); } else {
futures.push(disk.rename_part(src_bucket, src_object, dst_bucket, dst_object, meta.clone())) Err(Error::new(DiskError::DiskNotFound))
}
})
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -487,15 +491,15 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for (i, disk) in disks.iter().enumerate() { for (i, disk) in disks.iter().enumerate() {
if disk.is_none() {
errors.push(Some(Error::new(DiskError::DiskNotFound)));
continue;
}
let disk = disk.as_ref().unwrap();
let mut file_info = files[i].clone(); let mut file_info = files[i].clone();
file_info.erasure.index = i + 1; file_info.erasure.index = i + 1;
futures.push(disk.write_metadata(org_bucket, bucket, prefix, file_info)); futures.push(async move {
if let Some(disk) = disk {
disk.write_metadata(org_bucket, bucket, prefix, file_info).await
} else {
Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -961,25 +965,22 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() {
ress.push(FileInfo::default());
errors.push(Some(Error::new(DiskError::DiskNotFound)));
continue;
}
let disk = disk.as_ref().unwrap();
let opts = ReadOptions { read_data, healing }; let opts = ReadOptions { read_data, healing };
futures.push(async move { futures.push(async move {
if version_id.is_empty() { if let Some(disk) = disk {
match disk.read_xl(bucket, object, read_data).await { if version_id.is_empty() {
Ok(info) => { match disk.read_xl(bucket, object, read_data).await {
let fi = file_info_from_raw(info, bucket, object, read_data).await?; Ok(info) => {
Ok(fi) let fi = file_info_from_raw(info, bucket, object, read_data).await?;
Ok(fi)
}
Err(err) => Err(err),
} }
Err(err) => Err(err), } else {
disk.read_version(org_bucket, bucket, object, version_id, &opts).await
} }
} else { } else {
disk.read_version(org_bucket, bucket, object, version_id, &opts).await Err(Error::new(DiskError::DiskNotFound))
} }
}) })
} }
@@ -1023,14 +1024,13 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { futures.push(async move {
ress.push(None); if let Some(disk) = disk {
errors.push(Some(Error::new(DiskError::DiskNotFound))); disk.read_xl(bucket, object, read_data).await
continue; } else {
} Err(Error::new(DiskError::DiskNotFound))
}
let disk = disk.as_ref().unwrap(); });
futures.push(disk.read_xl(bucket, object, read_data));
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -1149,15 +1149,14 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() {
ress.push(None);
errors.push(Some(Error::new(DiskError::DiskNotFound)));
continue;
}
let disk = disk.as_ref().unwrap();
let req = req.clone(); let req = req.clone();
futures.push(disk.read_multiple(req)); futures.push(async move {
if let Some(disk) = disk {
disk.read_multiple(req).await
} else {
Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -1325,17 +1324,14 @@ impl SetDisks {
let mut ress = Vec::new(); let mut ress = Vec::new();
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() {
ress.push(None);
errs.push(Some(Error::new(DiskError::DiskNotFound)));
continue;
}
let disk = disk.as_ref().unwrap();
let opts = opts.clone(); let opts = opts.clone();
// let mut wr = &mut wr; futures.push(async move {
futures.push(disk.walk_dir(opts)); if let Some(disk) = disk {
// tokio::spawn(async move { disk.walk_dir(opts, wr).await }); disk.walk_dir(opts).await
} else {
Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -1375,20 +1371,18 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() {
errors.push(Some(Error::new(DiskError::DiskNotFound)));
continue;
}
let disk = disk.as_ref().unwrap();
let file_path = file_path.clone(); let file_path = file_path.clone();
let meta_file_path = format!("{}.meta", file_path); let meta_file_path = format!("{}.meta", file_path);
futures.push(async move { futures.push(async move {
disk.delete(RUSTFS_META_MULTIPART_BUCKET, &file_path, DeleteOptions::default()) if let Some(disk) = disk {
.await?; disk.delete(RUSTFS_META_MULTIPART_BUCKET, &file_path, DeleteOptions::default())
disk.delete(RUSTFS_META_MULTIPART_BUCKET, &meta_file_path, DeleteOptions::default()) .await?;
.await disk.delete(RUSTFS_META_MULTIPART_BUCKET, &meta_file_path, DeleteOptions::default())
.await
} else {
Err(Error::new(DiskError::DiskNotFound))
}
}); });
} }
@@ -1419,13 +1413,15 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { let file_path = file_path.clone();
errors.push(Some(Error::new(DiskError::DiskNotFound))); futures.push(async move {
continue; if let Some(disk) = disk {
} disk.delete(RUSTFS_META_MULTIPART_BUCKET, &file_path, DeleteOptions::default())
.await
let disk = disk.as_ref().unwrap(); } else {
futures.push(disk.delete(RUSTFS_META_MULTIPART_BUCKET, &file_path, DeleteOptions::default())); Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -1453,20 +1449,21 @@ impl SetDisks {
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { futures.push(async move {
errors.push(Some(Error::new(DiskError::DiskNotFound))); if let Some(disk) = disk {
continue; disk.delete(
} bucket,
prefix,
let disk = disk.as_ref().unwrap(); DeleteOptions {
futures.push(disk.delete( recursive: true,
bucket, ..Default::default()
prefix, },
DeleteOptions { )
recursive: true, .await
..Default::default() } else {
}, Err(Error::new(DiskError::DiskNotFound))
)); }
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -1626,9 +1623,9 @@ impl SetDisks {
let fi = Self::pick_valid_fileinfo(&parts_metadata, mot_time, etag, read_quorum as usize)?; let fi = Self::pick_valid_fileinfo(&parts_metadata, mot_time, etag, read_quorum as usize)?;
// debug!("get_object_fileinfo pick fi {:?}", &fi); // debug!("get_object_fileinfo pick fi {:?}", &fi);
let online_disks: Vec<Option<DiskStore>> = op_online_disks.iter().filter(|v| v.is_some()).cloned().collect(); // let online_disks: Vec<Option<DiskStore>> = op_online_disks.iter().filter(|v| v.is_some()).cloned().collect();
Ok((fi, parts_metadata, online_disks)) Ok((fi, parts_metadata, op_online_disks))
} }
#[allow(clippy::too_many_arguments)] #[allow(clippy::too_many_arguments)]
@@ -1763,12 +1760,14 @@ impl SetDisks {
let mut errs = Vec::with_capacity(disks.len()); let mut errs = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { let fi = fi.clone();
errs.push(Some(Error::new(DiskError::DiskNotFound))); futures.push(async move {
continue; if let Some(disk) = disk {
} disk.update_metadata(bucket, object, fi, opts).await
let disk = disk.as_ref().unwrap(); } else {
futures.push(disk.update_metadata(bucket, object, fi.clone(), opts)) Err(Error::new(DiskError::DiskNotFound))
}
})
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -3722,12 +3721,14 @@ impl StorageAPI for SetDisks {
// let mut errors = Vec::with_capacity(disks.len()); // let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { let vers = vers.clone();
continue; futures.push(async move {
} if let Some(disk) = disk {
disk.delete_versions(bucket, vers, DeleteOptions::default()).await
let disk = disk.as_ref().unwrap(); } else {
futures.push(disk.delete_versions(bucket, vers.clone(), DeleteOptions::default())); Err(Error::new(DiskError::DiskNotFound))
}
});
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -3960,19 +3961,16 @@ impl StorageAPI for SetDisks {
let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size);
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { if let Some(disk) = disk {
// let writer = disk.append_file(RUSTFS_META_TMP_BUCKET, &tmp_part_path).await?;
let filewriter = disk
.create_file("", RUSTFS_META_TMP_BUCKET, &tmp_part_path, data.content_length)
.await?;
let writer = new_bitrot_filewriter(filewriter, DEFAULT_BITROT_ALGO, erasure.shard_size(erasure.block_size));
writers.push(Some(writer));
} else {
writers.push(None); writers.push(None);
continue;
} }
let disk = disk.as_ref().unwrap().clone();
// let writer = disk.append_file(RUSTFS_META_TMP_BUCKET, &tmp_part_path).await?;
let filewriter = disk
.create_file("", RUSTFS_META_TMP_BUCKET, &tmp_part_path, data.content_length)
.await?;
let writer = new_bitrot_filewriter(filewriter, DEFAULT_BITROT_ALGO, erasure.shard_size(erasure.block_size));
writers.push(Some(writer));
} }
let mut erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); let mut erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size);
+17 -3
View File
@@ -11,7 +11,7 @@ use uuid::Uuid;
use crate::{ use crate::{
disk::{ disk::{
error::DiskError, error::{is_unformatted_disk, DiskError},
format::{DistributionAlgoVersion, FormatV3}, format::{DistributionAlgoVersion, FormatV3},
new_disk, DiskAPI, DiskInfo, DiskOption, DiskStore, new_disk, DiskAPI, DiskInfo, DiskOption, DiskStore,
}, },
@@ -34,8 +34,8 @@ use crate::{
use crate::heal::heal_ops::HealSequence; use crate::heal::heal_ops::HealSequence;
use tokio::time::Duration; use tokio::time::Duration;
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
use tracing::info;
use tracing::warn; use tracing::warn;
use tracing::{error, info};
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct Sets { pub struct Sets {
@@ -56,6 +56,7 @@ pub struct Sets {
} }
impl Sets { impl Sets {
#[tracing::instrument(level = "debug", skip(disks, endpoints, fm, pool_idx, partiy_count))]
pub async fn new( pub async fn new(
disks: Vec<Option<DiskStore>>, disks: Vec<Option<DiskStore>>,
endpoints: &PoolEndpoints, endpoints: &PoolEndpoints,
@@ -120,7 +121,20 @@ impl Sets {
disk = local_disk; disk = local_disk;
} }
if let Some(_disk_id) = disk.as_ref().unwrap().get_disk_id().await? { let has_disk_id = match disk.as_ref().unwrap().get_disk_id().await {
Ok(res) => res,
Err(err) => {
if is_unformatted_disk(&err) {
error!("get_disk_id err {:?}", err);
} else {
warn!("get_disk_id err {:?}", err);
}
None
}
};
if let Some(_disk_id) = has_disk_id {
set_drive.push(disk); set_drive.push(disk);
} else { } else {
warn!("sets new set_drive {}-{} get_disk_id is none", i, j); warn!("sets new set_drive {}-{} get_disk_id is none", i, j);
+1
View File
@@ -98,6 +98,7 @@ pub struct ECStore {
impl ECStore { impl ECStore {
#[allow(clippy::new_ret_no_self)] #[allow(clippy::new_ret_no_self)]
#[tracing::instrument(level = "debug", skip(endpoint_pools))]
pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<Arc<Self>> { pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<Arc<Self>> {
// let layouts = DisksLayout::from_volumes(endpoints.as_slice())?; // let layouts = DisksLayout::from_volumes(endpoints.as_slice())?;
+7 -7
View File
@@ -191,13 +191,13 @@ pub async fn load_format_erasure_all(disks: &[Option<DiskStore>], heal: bool) ->
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() { for disk in disks.iter() {
if disk.is_none() { futures.push(async move {
datas.push(None); if let Some(disk) = disk {
errors.push(Some(Error::new(DiskError::DiskNotFound))); load_format_erasure(disk, heal).await
} } else {
Err(Error::new(DiskError::DiskNotFound))
let disk = disk.as_ref().unwrap(); }
futures.push(load_format_erasure(disk, heal)); });
} }
let results = join_all(futures).await; let results = join_all(futures).await;
+6 -2
View File
@@ -209,10 +209,14 @@ async fn run(opt: config::Opt) -> Result<()> {
// init store // init store
let store = ECStore::new(server_address.clone(), endpoint_pools.clone()) let store = ECStore::new(server_address.clone(), endpoint_pools.clone())
.await .await
.map_err(|err| Error::from_string(err.to_string()))?; .map_err(|err| {
error!("ECStore::new {:?}", &err);
panic!("{}", err);
Error::from_string(err.to_string())
})?;
ECStore::init(store.clone()).await.map_err(|err| { ECStore::init(store.clone()).await.map_err(|err| {
error!("init faild {:?}", &err); error!("ECStore init faild {:?}", &err);
Error::from_string(err.to_string()) Error::from_string(err.to_string())
})?; })?;
warn!(" init store success!"); warn!(" init store success!");
+2 -2
View File
@@ -18,8 +18,8 @@ fi
export RUSTFS_STORAGE_CLASS_INLINE_BLOCK="512 KB" export RUSTFS_STORAGE_CLASS_INLINE_BLOCK="512 KB"
# DATA_DIR_ARG="./target/volume/test{0...4}" DATA_DIR_ARG="./target/volume/test{0...4}"
DATA_DIR_ARG="./target/volume/test" # DATA_DIR_ARG="./target/volume/test"
if [ -n "$1" ]; then if [ -n "$1" ]; then
DATA_DIR_ARG="$1" DATA_DIR_ARG="$1"