diff --git a/Cargo.toml b/Cargo.toml index 59c48cc14..684dfb53b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -190,7 +190,7 @@ inherits = "dev" [profile.release] opt-level = 3 -lto = "fat" +lto = "thin" codegen-units = 1 panic = "abort" # Optional, remove the panic expansion code strip = true # strip symbol information to reduce binary size diff --git a/ecstore/src/cache_value/metacache_set.rs b/ecstore/src/cache_value/metacache_set.rs index e716a3051..c3e75d093 100644 --- a/ecstore/src/cache_value/metacache_set.rs +++ b/ecstore/src/cache_value/metacache_set.rs @@ -7,7 +7,7 @@ use common::error::{Error, Result}; use futures::future::join_all; use std::{future::Future, pin::Pin, sync::Arc}; use tokio::{spawn, sync::broadcast::Receiver as B_Receiver}; -use tracing::{error, info}; +use tracing::{error, info, warn}; pub type AgreedFn = Box Pin + Send>> + Send + 'static>; pub type PartialFn = Box]) -> Pin + Send>> + Send + 'static>; @@ -205,7 +205,7 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - continue; } // If exact match, we agree. - if let Ok((_, true)) = current.matches(&entry, true) { + if let (_, true) = current.matches(Some(&entry), true) { top_entries[i] = Some(entry); agree += 1; @@ -224,7 +224,12 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - // We got a new, better current. // Clear existing entries. top_entries = vec![None; top_entries.len()]; - agree += 1; + + for item in top_entries.iter_mut().take(i) { + *item = None; + } + + agree = 1; top_entries[i] = Some(entry.clone()); current = entry; } @@ -265,6 +270,7 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - } } + warn!("list_path_raw: all at eof or error"); break; } @@ -272,6 +278,8 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - for r in readers.iter_mut() { let _ = r.skip(1).await; } + + warn!("list_path_raw: agree == readers.len() {} ", ¤t.name); if let Some(agreed_fn) = opts.agreed.as_ref() { agreed_fn(current).await; } @@ -285,6 +293,7 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - } } + warn!("list_path_raw: {} entries", top_entries.len()); if let Some(partial_fn) = opts.partial.as_ref() { partial_fn(MetaCacheEntries(top_entries), &errs).await; } diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index ddefd1bed..cda289c2b 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -786,46 +786,62 @@ impl MetaCacheEntry { fm.into_file_info_versions(bucket, self.name.as_str(), false) } - pub fn matches(&self, other: &MetaCacheEntry, strict: bool) -> Result<(Option, bool)> { + pub fn matches(&self, other: Option<&MetaCacheEntry>, strict: bool) -> (Option, bool) { + if other.is_none() { + return (None, false); + } + + let other = other.unwrap(); + let mut prefer = None; if self.name != other.name { if self.name < other.name { - return Ok((Some(self.clone()), false)); + return (Some(self.clone()), false); } - return Ok((Some(other.clone()), false)); + return (Some(other.clone()), false); } if other.is_dir() || self.is_dir() { if self.is_dir() { - return Ok((Some(self.clone()), other.is_dir())); + return (Some(self.clone()), other.is_dir() == self.is_dir()); } - return Ok((Some(other.clone()), other.is_dir() == self.is_dir())); + return (Some(other.clone()), other.is_dir() == self.is_dir()); } let self_vers = match &self.cached { Some(file_meta) => file_meta.clone(), - None => FileMeta::load(&self.metadata)?, + None => match FileMeta::load(&self.metadata) { + Ok(meta) => meta, + Err(_) => { + return (None, false); + } + }, }; let other_vers = match &other.cached { Some(file_meta) => file_meta.clone(), - None => FileMeta::load(&other.metadata)?, + None => match FileMeta::load(&other.metadata) { + Ok(meta) => meta, + Err(_) => { + return (None, false); + } + }, }; if self_vers.versions.len() != other_vers.versions.len() { match self_vers.lastest_mod_time().cmp(&other_vers.lastest_mod_time()) { Ordering::Greater => { - return Ok((Some(self.clone()), false)); + return (Some(self.clone()), false); } Ordering::Less => { - return Ok((Some(self.clone()), false)); + return (Some(other.clone()), false); } _ => {} } if self_vers.versions.len() > other_vers.versions.len() { - return Ok((Some(self.clone()), false)); + return (Some(self.clone()), false); } - return Ok((Some(self.clone()), false)); + return (Some(other.clone()), false); } for (s_version, o_version) in self_vers.versions.iter().zip(other_vers.versions.iter()) { @@ -853,14 +869,14 @@ impl MetaCacheEntry { } if prefer.is_some() { - return Ok((prefer, false)); + return (prefer, false); } if s_version.header.sorts_before(&o_version.header) { - return Ok((Some(self.clone()), false)); + return (Some(self.clone()), false); } - return Ok((Some(other.clone()), false)); + return (Some(other.clone()), false); } } @@ -868,7 +884,7 @@ impl MetaCacheEntry { prefer = Some(self.clone()); } - Ok((prefer, true)) + (prefer, true) } pub fn xl_meta(&mut self) -> Result { @@ -900,9 +916,10 @@ impl MetaCacheEntries { pub fn as_ref(&self) -> &[Option] { &self.0 } - pub fn resolve(&self, mut params: MetadataResolutionParams) -> Result> { + pub fn resolve(&self, mut params: MetadataResolutionParams) -> Option { if self.0.is_empty() { - return Ok(None); + warn!("decommission_pool: entries resolve empty"); + return None; } let mut dir_exists = 0; @@ -913,76 +930,104 @@ impl MetaCacheEntries { let mut objs_valid = 0; for entry in self.0.iter().flatten() { + let mut entry = entry.clone(); + + warn!("decommission_pool: entries resolve entry {:?}", entry.name); if entry.name.is_empty() { continue; } if entry.is_dir() { dir_exists += 1; selected = Some(entry.clone()); + warn!("decommission_pool: entries resolve entry dir {:?}", entry.name); continue; } + let xl = match entry.xl_meta() { + Ok(xl) => xl, + Err(e) => { + warn!("decommission_pool: entries resolve entry xl_meta {:?}", e); + continue; + } + }; + objs_valid += 1; - match &entry.cached { - Some(file_meta) => { - params.candidates.push(file_meta.versions.clone()); - } - None => { - params.candidates.push(FileMeta::load(&entry.metadata)?.versions); - } - } + params.candidates.push(xl.versions.clone()); if selected.is_none() { selected = Some(entry.clone()); objs_agree = 1; + warn!("decommission_pool: entries resolve entry selected {:?}", entry.name); continue; } - if let (Some(prefer), true) = entry.matches(selected.as_ref().unwrap(), params.strict)? { - selected = Some(prefer); + if let (prefer, true) = entry.matches(selected.as_ref(), params.strict) { + selected = prefer; objs_agree += 1; + warn!("decommission_pool: entries resolve entry prefer {:?}", entry.name); continue; } } - // Return dir entries, if enough... - if selected.is_some() && selected.as_ref().unwrap().is_dir() && dir_exists >= params.dir_quorum { - return Ok(selected); + let Some(selected) = selected else { + warn!("decommission_pool: entries resolve entry no selected"); + return None; + }; + + if selected.is_dir() && dir_exists >= params.dir_quorum { + warn!("decommission_pool: entries resolve entry dir selected {:?}", selected.name); + return Some(selected); } + // If we would never be able to reach read quorum. if objs_valid < params.obj_quorum { - return Ok(None); + warn!( + "decommission_pool: entries resolve entry not enough objects {} < {}", + objs_valid, params.obj_quorum + ); + return None; } - // If all objects agree. - if selected.is_some() && objs_agree == objs_valid { - return Ok(selected); + + if objs_agree == objs_valid { + warn!("decommission_pool: entries resolve entry all agree {} == {}", objs_agree, objs_valid); + return Some(selected); } - // If cached is nil we shall skip the entry. - if selected.is_none() || (selected.is_some() && selected.as_ref().unwrap().cached.is_none()) { - return Ok(None); + + let Some(cached) = selected.cached else { + warn!("decommission_pool: entries resolve entry no cached"); + return None; + }; + + let versions = merge_file_meta_versions(params.obj_quorum, params.strict, params.requested_versions, ¶ms.candidates); + if versions.is_empty() { + warn!("decommission_pool: entries resolve entry no versions"); + return None; } + + let metadata = match cached.marshal_msg() { + Ok(meta) => meta, + Err(e) => { + warn!("decommission_pool: entries resolve entry marshal_msg {:?}", e); + return None; + } + }; + // Merge if we have disagreement. // Create a new merged result. - selected = Some(MetaCacheEntry { - name: selected.as_ref().unwrap().name.clone(), + let new_selected = MetaCacheEntry { + name: selected.name.clone(), cached: Some(FileMeta { - meta_ver: selected.as_ref().unwrap().cached.as_ref().unwrap().meta_ver, + meta_ver: cached.meta_ver, + versions, ..Default::default() }), reusable: true, - ..Default::default() - }); + metadata, + }; - selected.as_mut().unwrap().cached.as_mut().unwrap().versions = - merge_file_meta_versions(params.obj_quorum, params.strict, params.requested_versions, ¶ms.candidates); - if selected.as_ref().unwrap().cached.as_ref().unwrap().versions.is_empty() { - return Ok(None); - } - - selected.as_mut().unwrap().metadata = selected.as_ref().unwrap().cached.as_ref().unwrap().marshal_msg()?; - - Ok(selected) + warn!("decommission_pool: entries resolve entry selected {:?}", new_selected.name); + Some(new_selected) } pub fn first_found(&self) -> (Option, usize) { diff --git a/ecstore/src/file_meta.rs b/ecstore/src/file_meta.rs index e4da3882a..1234c07f4 100644 --- a/ecstore/src/file_meta.rs +++ b/ecstore/src/file_meta.rs @@ -916,11 +916,20 @@ impl FileMetaVersionHeader { } pub fn matches_not_strict(&self, o: &FileMetaVersionHeader) -> bool { + let mut ok = self.version_id == o.version_id && self.version_type == o.version_type && self.matches_ec(o); if self.version_id.is_none() { - return self.version_id == o.version_id && self.version_type == o.version_type && self.mod_time == o.mod_time; + ok = ok && self.mod_time == o.mod_time; } - self.version_id == o.version_id && self.version_type == o.version_type + ok + } + + pub fn matches_ec(&self, o: &FileMetaVersionHeader) -> bool { + if self.has_ec() && o.has_ec() { + return self.ec_n == o.ec_n && self.ec_m == o.ec_m; + } + + true } pub fn free_version(&self) -> bool { diff --git a/ecstore/src/heal/data_scanner.rs b/ecstore/src/heal/data_scanner.rs index b1b110a54..c806267f2 100644 --- a/ecstore/src/heal/data_scanner.rs +++ b/ecstore/src/heal/data_scanner.rs @@ -867,7 +867,7 @@ impl FolderScanner { // return; // } let entry = match entries.resolve(resolver_partial) { - Ok(Some(entry)) => entry, + Some(entry) => entry, _ => match entries.first_found() { (Some(entry), _) => entry, _ => return, diff --git a/ecstore/src/pools.rs b/ecstore/src/pools.rs index f6f996941..c151b0627 100644 --- a/ecstore/src/pools.rs +++ b/ecstore/src/pools.rs @@ -1330,22 +1330,17 @@ impl SetDisks { partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { let resolver = resolver.clone(); let cb_func = cb_func.clone(); - match entries.resolve(resolver) { - Ok(Some(entry)) => { + Some(entry) => { warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name); Box::pin(async move { cb_func(entry).await; }) } - Ok(None) => { + None => { warn!("decommission_pool: list_objects_to_decommission get none"); Box::pin(async {}) } - Err(err) => { - error!("decommission_pool: list_objects_to_decommission get err {:?}", &err); - Box::pin(async {}) - } } })), ..Default::default() diff --git a/ecstore/src/rebalance.rs b/ecstore/src/rebalance.rs index 62c5d696f..d6ec57113 100644 --- a/ecstore/src/rebalance.rs +++ b/ecstore/src/rebalance.rs @@ -1075,18 +1075,14 @@ impl SetDisks { let cb = cb.clone(); match entries.resolve(resolver) { - Ok(Some(entry)) => { + Some(entry) => { warn!("rebalance: list_objects_to_decommission get {}", &entry.name); Box::pin(async move { cb(entry).await }) } - Ok(None) => { + None => { warn!("rebalance: list_objects_to_decommission get none"); Box::pin(async {}) } - Err(err) => { - error!("rebalance: list_objects_to_decommission get err {:?}", &err); - Box::pin(async {}) - } } })), ..Default::default() diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index bd55ab39f..da3046a7b 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -2052,7 +2052,7 @@ impl SetDisks { let bucket_partial = bucket_partial.clone(); async move { let entry = match entries.resolve(resolver_partial) { - Ok(Some(entry)) => entry, + Some(entry) => entry, _ => match entries.first_found() { (Some(entry), _) => entry, _ => return, @@ -3505,7 +3505,7 @@ impl SetDisks { let heal_entry = heal_entry.clone(); let resolver = resolver.clone(); async move { - let entry = if let Ok(Some(entry)) = entries.resolve(resolver) { + let entry = if let Some(entry) = entries.resolve(resolver) { entry } else if let (Some(entry), _) = entries.first_found() { entry diff --git a/ecstore/src/store_list_objects.rs b/ecstore/src/store_list_objects.rs index 8b9118de3..445c56f0c 100644 --- a/ecstore/src/store_list_objects.rs +++ b/ecstore/src/store_list_objects.rs @@ -778,7 +778,7 @@ impl ECStore { let value = tx2.clone(); let resolver = resolver.clone(); async move { - if let Ok(Some(entry)) = entries.resolve(resolver) { + if let Some(entry) = entries.resolve(resolver) { if let Err(err) = value.send(entry).await { error!("list_path send fail {:?}", err); } @@ -1296,7 +1296,7 @@ impl SetDisks { let value = tx2.clone(); let resolver = resolver.clone(); async move { - if let Ok(Some(entry)) = entries.resolve(resolver) { + if let Some(entry) = entries.resolve(resolver) { if let Err(err) = value.send(entry).await { error!("list_path send fail {:?}", err); } diff --git a/scripts/run.sh b/scripts/run.sh index 690ac8511..352568ecc 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -33,7 +33,7 @@ export RUSTFS_CONSOLE_ENABLE=true export RUSTFS_CONSOLE_ADDRESS=":9002" # export RUSTFS_SERVER_DOMAINS="localhost:9000" # HTTPS 证书目录 - export RUSTFS_TLS_PATH="./deploy/certs" +# export RUSTFS_TLS_PATH="./deploy/certs" # 具体路径修改为配置文件真实路径,obs.example.toml 仅供参考 其中`RUSTFS_OBS_CONFIG` 和下面变量二选一 export RUSTFS_OBS_CONFIG="./deploy/config/obs.example.toml"