mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-13 00:26:53 +00:00
merge main
This commit is contained in:
+84
-9
@@ -6,7 +6,10 @@ use std::{
|
||||
time::Duration,
|
||||
};
|
||||
|
||||
use crate::config::error::is_not_found;
|
||||
use crate::global::GLOBAL_MRFState;
|
||||
use crate::heal::heal_ops::{HealEntryFn, HealSequence};
|
||||
use crate::heal::mrf::PartialOperation;
|
||||
use crate::{
|
||||
bitrot::{bitrot_verify, close_bitrot_writers, new_bitrot_filereader, new_bitrot_filewriter, BitrotFileWriter},
|
||||
cache_value::metacache_set::{list_path_raw, ListPathRawOptions},
|
||||
@@ -46,7 +49,7 @@ use crate::{
|
||||
store_init::{load_format_erasure, ErasureError},
|
||||
utils::{
|
||||
self,
|
||||
crypto::{base64_decode, base64_encode, hex, sha256},
|
||||
crypto::{base64_decode, base64_encode, hex},
|
||||
path::{encode_dir_object, has_suffix, SLASH_SEPARATOR},
|
||||
},
|
||||
xhttp,
|
||||
@@ -56,6 +59,7 @@ use crate::{
|
||||
heal::data_scanner::{globalHealConfig, HEAL_DELETE_DANGLING},
|
||||
store_api::ListObjectVersionsInfo,
|
||||
};
|
||||
use chrono::Utc;
|
||||
use futures::future::join_all;
|
||||
use glob::Pattern;
|
||||
use http::HeaderMap;
|
||||
@@ -619,7 +623,9 @@ impl SetDisks {
|
||||
|
||||
fn get_multipart_sha_dir(bucket: &str, object: &str) -> String {
|
||||
let path = format!("{}/{}", bucket, object);
|
||||
hex(sha256(path.as_bytes()).as_ref())
|
||||
let mut hasher = Sha256::new();
|
||||
hasher.update(path);
|
||||
hex(hasher.finalize())
|
||||
}
|
||||
|
||||
fn common_parity(parities: &[i32], default_parity_count: i32) -> i32 {
|
||||
@@ -1318,6 +1324,7 @@ impl SetDisks {
|
||||
for (i, opdisk) in disks.iter().enumerate() {
|
||||
if let Some(disk) = opdisk {
|
||||
if disk.is_online().await && disk.get_disk_location().set_idx.is_some() {
|
||||
info!("Disk {:?} is online", disk);
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -1325,6 +1332,7 @@ impl SetDisks {
|
||||
}
|
||||
|
||||
if let Some(endpoint) = self.set_endpoints.get(i) {
|
||||
info!("will renew disk, opdisk: {:?}", opdisk);
|
||||
self.renew_disk(endpoint).await;
|
||||
}
|
||||
}
|
||||
@@ -1338,12 +1346,20 @@ impl SetDisks {
|
||||
Err(e) => {
|
||||
warn!("connect_endpoint err {:?}", &e);
|
||||
if ep.is_local && DiskError::UnformattedDisk.is(&e) {
|
||||
// TODO: pushHealLocalDisks
|
||||
GLOBAL_BackgroundHealState.push_heal_local_disks(&[ep.clone()]).await;
|
||||
}
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
if new_disk.is_local() {
|
||||
if let Some(h) = new_disk.healing().await {
|
||||
if !h.finished {
|
||||
GLOBAL_BackgroundHealState.push_heal_local_disks(&[new_disk.endpoint()]).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let (set_idx, disk_idx) = match self.find_disk_index(&fm) {
|
||||
Ok(res) => res,
|
||||
Err(e) => {
|
||||
@@ -1706,6 +1722,19 @@ impl SetDisks {
|
||||
let (op_online_disks, mot_time, etag) = Self::list_online_disks(&disks, &parts_metadata, &errs, read_quorum as usize);
|
||||
|
||||
let fi = Self::pick_valid_fileinfo(&parts_metadata, mot_time, etag, read_quorum as usize)?;
|
||||
if errs.iter().any(|err| err.is_some()) {
|
||||
GLOBAL_MRFState
|
||||
.add_partial(PartialOperation {
|
||||
bucket: fi.volume.to_string(),
|
||||
object: fi.name.to_string(),
|
||||
queued: Utc::now(),
|
||||
version_id: fi.version_id.map(|v| v.to_string()),
|
||||
set_index: self.set_index,
|
||||
pool_index: self.pool_index,
|
||||
..Default::default()
|
||||
})
|
||||
.await;
|
||||
}
|
||||
// debug!("get_object_fileinfo pick fi {:?}", &fi);
|
||||
|
||||
// let online_disks: Vec<Option<DiskStore>> = op_online_disks.iter().filter(|v| v.is_some()).cloned().collect();
|
||||
@@ -1724,6 +1753,8 @@ impl SetDisks {
|
||||
fi: FileInfo,
|
||||
files: Vec<FileInfo>,
|
||||
disks: &[Option<DiskStore>],
|
||||
set_index: usize,
|
||||
pool_index: usize,
|
||||
) -> Result<()> {
|
||||
let (disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, &files, &fi);
|
||||
|
||||
@@ -1809,10 +1840,34 @@ impl SetDisks {
|
||||
// "read part {} part_offset {},part_length {},part_size {} ",
|
||||
// part_number, part_offset, part_length, part_size
|
||||
// );
|
||||
let _n = erasure
|
||||
let (written, mut err) = erasure
|
||||
.decode(writer, readers, part_offset as usize, part_length, part_size)
|
||||
.await?;
|
||||
|
||||
.await;
|
||||
if let Some(e) = err.as_ref() {
|
||||
if written == part_length {
|
||||
match e.downcast_ref::<DiskError>() {
|
||||
Some(DiskError::FileNotFound) | Some(DiskError::FileCorrupt) => {
|
||||
GLOBAL_MRFState
|
||||
.add_partial(PartialOperation {
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
queued: Utc::now(),
|
||||
version_id: fi.version_id.map(|v| v.to_string()),
|
||||
set_index,
|
||||
pool_index,
|
||||
bitrot_scan: !is_not_found(e),
|
||||
..Default::default()
|
||||
})
|
||||
.await;
|
||||
err = None;
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(err) = err {
|
||||
return Err(err);
|
||||
}
|
||||
// debug!("ec decode {} writed size {}", part_number, n);
|
||||
|
||||
total_readed += part_length as i64;
|
||||
@@ -2941,7 +2996,6 @@ impl SetDisks {
|
||||
info!("ns_scanner start");
|
||||
let _ = join_all(futures).await;
|
||||
drop(buckets_results_tx);
|
||||
info!("1");
|
||||
let _ = task.await;
|
||||
info!("ns_scanner completed");
|
||||
Ok(())
|
||||
@@ -3070,6 +3124,7 @@ impl SetDisks {
|
||||
let mut ret_err = None;
|
||||
for bucket in buckets.iter() {
|
||||
if tracker.read().await.is_healed(bucket).await {
|
||||
info!("bucket{} was healed", bucket);
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -3190,6 +3245,7 @@ impl SetDisks {
|
||||
let bg_seq_clone = bg_seq.clone();
|
||||
let send_clone = send.clone();
|
||||
let heal_entry = Arc::new(move |bucket: String, entry: MetaCacheEntry| {
|
||||
info!("heal entry, bucket: {}, entry: {:?}", bucket, entry);
|
||||
let jt_clone = jt_clone.clone();
|
||||
let self_clone = self_clone.clone();
|
||||
let started = started_clone;
|
||||
@@ -3469,8 +3525,14 @@ impl ObjectIO for SetDisks {
|
||||
// let disks = disks.clone();
|
||||
let bucket = String::from(bucket);
|
||||
let object = String::from(object);
|
||||
let set_index = self.set_index.clone();
|
||||
let pool_index = self.pool_index.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = Self::get_object_with_fileinfo(&bucket, &object, offset, length, &mut wd, fi, files, &disks).await {
|
||||
if let Err(e) = Self::get_object_with_fileinfo(
|
||||
&bucket, &object, offset, length, &mut wd, fi, files, &disks, set_index, pool_index,
|
||||
)
|
||||
.await
|
||||
{
|
||||
error!("get_object_with_fileinfo err {:?}", e);
|
||||
};
|
||||
});
|
||||
@@ -4529,7 +4591,7 @@ impl StorageAPI for SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
let (online_disks, _, op_old_dir) = Self::rename_data(
|
||||
let (online_disks, versions, op_old_dir) = Self::rename_data(
|
||||
&shuffle_disks,
|
||||
RUSTFS_META_MULTIPART_BUCKET,
|
||||
&upload_id_path,
|
||||
@@ -4558,6 +4620,19 @@ impl StorageAPI for SetDisks {
|
||||
self.commit_rename_data_dir(&shuffle_disks, bucket, object, &old_dir.to_string(), write_quorum)
|
||||
.await?;
|
||||
}
|
||||
if let Some(versions) = versions {
|
||||
GLOBAL_MRFState
|
||||
.add_partial(PartialOperation {
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
queued: Utc::now(),
|
||||
versions,
|
||||
set_index: self.set_index,
|
||||
pool_index: self.pool_index,
|
||||
..Default::default()
|
||||
})
|
||||
.await;
|
||||
}
|
||||
|
||||
let _ = self.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user