mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 00:38:16 +00:00
Merge pull request #391 from rustfs/dada/fix-entry
feat: decom/rebalance
This commit is contained in:
@@ -34,7 +34,7 @@ docker compose -f docker-compose.yml up -d
|
||||
| service_name | 服务名称 | rustfs |
|
||||
| service_version | 服务版本 | 1.0.0 |
|
||||
| environment | 运行环境 | production |
|
||||
| meter_interval | 指标导出间隔 (秒) | 30 |
|
||||
| meter_interval | 指标导出间隔 (秒) | 30 |
|
||||
| sample_ratio | 采样率 | 1.0 |
|
||||
| use_stdout | 是否输出到控制台 | true/false |
|
||||
| logger_level | 日志级别 | info |
|
||||
|
||||
+1
-1
@@ -204,7 +204,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
|
||||
|
||||
@@ -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;
|
||||
|
||||
pub type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
|
||||
pub type PartialFn = Box<dyn Fn(MetaCacheEntries, &[Option<Error>]) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
|
||||
@@ -54,7 +54,6 @@ impl Clone for ListPathRawOptions {
|
||||
pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -> Result<()> {
|
||||
// println!("list_path_raw {},{}", &opts.bucket, &opts.path);
|
||||
if opts.disks.is_empty() {
|
||||
info!("list_path_raw 0 drives provided");
|
||||
return Err(Error::from_string("list_path_raw: 0 drives provided"));
|
||||
}
|
||||
|
||||
@@ -205,7 +204,7 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, 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;
|
||||
|
||||
@@ -214,17 +213,18 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
|
||||
// If only the name matches we didn't agree, but add it for resolution.
|
||||
if entry.name == current.name {
|
||||
top_entries[i] = Some(entry);
|
||||
|
||||
continue;
|
||||
}
|
||||
// We got different entries
|
||||
if entry.name > current.name {
|
||||
continue;
|
||||
}
|
||||
// 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;
|
||||
}
|
||||
@@ -272,6 +272,7 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
|
||||
for r in readers.iter_mut() {
|
||||
let _ = r.skip(1).await;
|
||||
}
|
||||
|
||||
if let Some(agreed_fn) = opts.agreed.as_ref() {
|
||||
agreed_fn(current).await;
|
||||
}
|
||||
|
||||
+96
-51
@@ -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<MetaCacheEntry>, bool)> {
|
||||
pub fn matches(&self, other: Option<&MetaCacheEntry>, strict: bool) -> (Option<MetaCacheEntry>, 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<FileMeta> {
|
||||
@@ -900,9 +916,10 @@ impl MetaCacheEntries {
|
||||
pub fn as_ref(&self) -> &[Option<MetaCacheEntry>] {
|
||||
&self.0
|
||||
}
|
||||
pub fn resolve(&self, mut params: MetadataResolutionParams) -> Result<Option<MetaCacheEntry>> {
|
||||
pub fn resolve(&self, mut params: MetadataResolutionParams) -> Option<MetaCacheEntry> {
|
||||
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<MetaCacheEntry>, usize) {
|
||||
|
||||
@@ -920,11 +920,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 {
|
||||
|
||||
@@ -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,
|
||||
|
||||
+26
-17
@@ -31,6 +31,7 @@ use std::io::{Cursor, Write};
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use time::{Duration, OffsetDateTime};
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::broadcast::Receiver as B_Receiver;
|
||||
use tracing::{error, info, warn};
|
||||
|
||||
@@ -910,7 +911,11 @@ impl ECStore {
|
||||
let wk = wk.clone();
|
||||
let set = set.clone();
|
||||
let rcfg = rcfg.clone();
|
||||
Box::pin(async move { this.decommission_entry(idx, entry, bucket, set, wk, rcfg).await })
|
||||
|
||||
Box::pin(async move {
|
||||
wk.take().await;
|
||||
this.decommission_entry(idx, entry, bucket, set, wk, rcfg).await
|
||||
})
|
||||
}
|
||||
});
|
||||
|
||||
@@ -918,6 +923,7 @@ impl ECStore {
|
||||
let mut rx = rx.resubscribe();
|
||||
let bi = bi.clone();
|
||||
let set_id = set_idx;
|
||||
let wk_clone = wk.clone();
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
if rx.try_recv().is_ok() {
|
||||
@@ -945,6 +951,8 @@ impl ECStore {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
wk_clone.give().await;
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1166,7 +1174,7 @@ impl ECStore {
|
||||
|
||||
#[tracing::instrument(skip(self, rd))]
|
||||
async fn decommission_object(self: Arc<Self>, pool_idx: usize, bucket: String, rd: GetObjectReader) -> Result<()> {
|
||||
warn!("decommission_object: {} {}", &bucket, &rd.object_info.name);
|
||||
warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name);
|
||||
let object_info = rd.object_info.clone();
|
||||
|
||||
// TODO: check : use size or actual_size ?
|
||||
@@ -1194,8 +1202,6 @@ impl ECStore {
|
||||
}
|
||||
};
|
||||
|
||||
// TODO: defer abort_multipart_upload
|
||||
|
||||
defer!(|| async {
|
||||
if let Err(err) = self
|
||||
.abort_multipart_upload(&bucket, &object_info.name, &res.upload_id, &ObjectOptions::default())
|
||||
@@ -1210,9 +1216,13 @@ impl ECStore {
|
||||
let mut reader = rd.stream;
|
||||
|
||||
for (i, part) in object_info.parts.iter().enumerate() {
|
||||
// 每次从reader中读取一个part上传
|
||||
let mut chunk = vec![0u8; part.size];
|
||||
|
||||
let mut data = PutObjReader::new(reader, part.size);
|
||||
reader.read_exact(&mut chunk).await?;
|
||||
|
||||
// 每次从reader中读取一个part上传
|
||||
let rd = Box::new(Cursor::new(chunk));
|
||||
let mut data = PutObjReader::new(rd, part.size);
|
||||
|
||||
let pi = match self
|
||||
.put_object_part(
|
||||
@@ -1230,18 +1240,17 @@ impl ECStore {
|
||||
{
|
||||
Ok(pi) => pi,
|
||||
Err(err) => {
|
||||
error!("decommission_object: put_object_part err {:?}", &err);
|
||||
error!("decommission_object: put_object_part {} err {:?}", i, &err);
|
||||
return Err(err);
|
||||
}
|
||||
};
|
||||
|
||||
warn!("decommission_object: put_object_part {} done {} {}", i, &bucket, &object_info.name);
|
||||
|
||||
parts[i] = CompletePart {
|
||||
part_num: pi.part_num,
|
||||
e_tag: pi.etag,
|
||||
};
|
||||
|
||||
// 把reader所有权拿回来?
|
||||
reader = data.stream;
|
||||
}
|
||||
|
||||
if let Err(err) = self
|
||||
@@ -1262,6 +1271,7 @@ impl ECStore {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
warn!("decommission_object: complete_multipart_upload done {} {}", &bucket, &object_info.name);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
@@ -1289,6 +1299,7 @@ impl ECStore {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
warn!("decommission_object: put_object done {} {}", &bucket, &object_info.name);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -1319,6 +1330,8 @@ impl SetDisks {
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let cb1 = cb_func.clone();
|
||||
|
||||
list_path_raw(
|
||||
rx,
|
||||
ListPathRawOptions {
|
||||
@@ -1327,25 +1340,21 @@ impl SetDisks {
|
||||
path: bucket_info.prefix.clone(),
|
||||
recursice: true,
|
||||
min_disks: listing_quorum,
|
||||
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))),
|
||||
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<Error>]| {
|
||||
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()
|
||||
|
||||
+25
-24
@@ -18,6 +18,7 @@ use common::defer;
|
||||
use common::error::{Error, Result};
|
||||
use http::HeaderMap;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::broadcast::{self, Receiver as B_Receiver};
|
||||
use tokio::time::{Duration, Instant};
|
||||
use tracing::{error, info, warn};
|
||||
@@ -370,7 +371,7 @@ impl ECStore {
|
||||
let rebalance_meta = self.rebalance_meta.read().await;
|
||||
if let Some(meta) = rebalance_meta.as_ref() {
|
||||
if let Some(pool_stat) = meta.pool_stats.get(pool_index) {
|
||||
if pool_stat.info.status != RebalStatus::Completed || !pool_stat.participating {
|
||||
if pool_stat.info.status == RebalStatus::Completed || !pool_stat.participating {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
@@ -389,10 +390,15 @@ impl ECStore {
|
||||
let mut rebalance_meta = self.rebalance_meta.write().await;
|
||||
if let Some(meta) = rebalance_meta.as_mut() {
|
||||
if let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) {
|
||||
if let Some(idx) = pool_stat.buckets.iter().position(|b| *b == bucket) {
|
||||
warn!("bucket_rebalance_done: buckets {:?}", &pool_stat.buckets);
|
||||
if let Some(idx) = pool_stat.buckets.iter().position(|b| b.as_str() == bucket.as_str()) {
|
||||
warn!("bucket_rebalance_done: bucket {} rebalanced", &bucket);
|
||||
pool_stat.buckets.remove(idx);
|
||||
pool_stat.rebalanced_buckets.push(bucket);
|
||||
|
||||
return Ok(());
|
||||
} else {
|
||||
warn!("bucket_rebalance_done: bucket {} not found", bucket);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -581,25 +587,26 @@ impl ECStore {
|
||||
}
|
||||
});
|
||||
|
||||
tracing::warn!("Pool {} rebalancing is started", pool_index + 1);
|
||||
warn!("Pool {} rebalancing is started", pool_index + 1);
|
||||
|
||||
while let Some(bucket) = self.next_rebal_bucket(pool_index).await? {
|
||||
tracing::info!("Rebalancing bucket: {}", bucket);
|
||||
warn!("Rebalancing bucket: start {}", bucket);
|
||||
|
||||
if let Err(err) = self.rebalance_bucket(rx.resubscribe(), bucket.clone(), pool_index).await {
|
||||
if err.to_string().contains("not initialized") {
|
||||
warn!("rebalance_bucket: rebalance not initialized, continue");
|
||||
continue;
|
||||
}
|
||||
tracing::error!("Error rebalancing bucket {}: {:?}", bucket, err);
|
||||
error!("Error rebalancing bucket {}: {:?}", bucket, err);
|
||||
done_tx.send(Err(err)).await.ok();
|
||||
break;
|
||||
}
|
||||
|
||||
warn!("Rebalance bucket: done {} ", bucket);
|
||||
self.bucket_rebalance_done(pool_index, bucket).await?;
|
||||
}
|
||||
|
||||
tracing::warn!("Pool {} rebalancing is done", pool_index + 1);
|
||||
warn!("Pool {} rebalancing is done", pool_index + 1);
|
||||
|
||||
done_tx.send(Ok(())).await.ok();
|
||||
save_task.await.ok();
|
||||
@@ -632,7 +639,6 @@ impl ECStore {
|
||||
false
|
||||
}
|
||||
|
||||
#[allow(unused_assignments)]
|
||||
#[tracing::instrument(skip(self, wk, set))]
|
||||
async fn rebalance_entry(
|
||||
&self,
|
||||
@@ -838,8 +844,6 @@ impl ECStore {
|
||||
}
|
||||
};
|
||||
|
||||
// TODO: defer abort_multipart_upload
|
||||
|
||||
defer!(|| async {
|
||||
if let Err(err) = self
|
||||
.abort_multipart_upload(&bucket, &object_info.name, &res.upload_id, &ObjectOptions::default())
|
||||
@@ -856,7 +860,13 @@ impl ECStore {
|
||||
for (i, part) in object_info.parts.iter().enumerate() {
|
||||
// 每次从reader中读取一个part上传
|
||||
|
||||
let mut data = PutObjReader::new(reader, part.size);
|
||||
let mut chunk = vec![0u8; part.size];
|
||||
|
||||
reader.read_exact(&mut chunk).await?;
|
||||
|
||||
// 每次从reader中读取一个part上传
|
||||
let rd = Box::new(Cursor::new(chunk));
|
||||
let mut data = PutObjReader::new(rd, part.size);
|
||||
|
||||
let pi = match self
|
||||
.put_object_part(
|
||||
@@ -883,9 +893,6 @@ impl ECStore {
|
||||
part_num: pi.part_num,
|
||||
e_tag: pi.etag,
|
||||
};
|
||||
|
||||
// 把reader所有权拿回来?
|
||||
reader = data.stream;
|
||||
}
|
||||
|
||||
if let Err(err) = self
|
||||
@@ -939,7 +946,7 @@ impl ECStore {
|
||||
#[tracing::instrument(skip(self, rx))]
|
||||
async fn rebalance_bucket(self: &Arc<Self>, rx: B_Receiver<bool>, bucket: String, pool_index: usize) -> Result<()> {
|
||||
// Placeholder for actual bucket rebalance logic
|
||||
tracing::info!("Rebalancing bucket {} in pool {}", bucket, pool_index);
|
||||
warn!("Rebalancing bucket {} in pool {}", bucket, pool_index);
|
||||
|
||||
// TODO: other config
|
||||
// if bucket != RUSTFS_META_BUCKET{
|
||||
@@ -977,14 +984,12 @@ impl ECStore {
|
||||
let bucket = bucket.clone();
|
||||
let wk = wk.clone();
|
||||
tokio::spawn(async move {
|
||||
defer!(|| async {
|
||||
wk.clone().give().await;
|
||||
});
|
||||
if let Err(err) = set.list_objects_to_rebalance(rx, bucket, rebalance_entry).await {
|
||||
error!("Rebalance worker {} error: {}", set_idx, err);
|
||||
} else {
|
||||
info!("Rebalance worker {} done", set_idx);
|
||||
}
|
||||
wk.clone().give().await;
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1068,25 +1073,21 @@ impl SetDisks {
|
||||
bucket: bucket.clone(),
|
||||
recursice: true,
|
||||
min_disks: listing_quorum,
|
||||
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry.clone())))),
|
||||
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))),
|
||||
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<Error>]| {
|
||||
// let cb = cb.clone();
|
||||
let resolver = resolver.clone();
|
||||
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()
|
||||
|
||||
@@ -2063,7 +2063,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,
|
||||
@@ -3516,7 +3516,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
|
||||
|
||||
+30
-19
@@ -524,7 +524,9 @@ impl ECStore {
|
||||
|
||||
// TODO: 并发
|
||||
for (idx, pool) in self.pools.iter().enumerate() {
|
||||
// TODO: IsSuspended
|
||||
if self.is_suspended(idx).await || self.is_pool_rebalancing(idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
n_sets[idx] = pool.set_count;
|
||||
|
||||
@@ -714,16 +716,14 @@ impl ECStore {
|
||||
let mut def_pool = PoolObjInfo::default();
|
||||
let mut has_def_pool = false;
|
||||
|
||||
let pool_meta = self.pool_meta.read().await;
|
||||
for pinfo in ress.iter() {
|
||||
if opts.skip_decommissioned && pool_meta.is_suspended(pinfo.index) {
|
||||
if opts.skip_decommissioned && self.is_suspended(pinfo.index).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
// TODO:SkipRebalancing
|
||||
// if opts.SkipRebalancing && z.IsPoolRebalancing(pinfo.Index) {
|
||||
// continue
|
||||
// }
|
||||
if opts.skip_rebalancing && self.is_pool_rebalancing(pinfo.index).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
if pinfo.err.is_none() {
|
||||
return Ok((pinfo.clone(), self.pools_with_object(&ress, opts).await));
|
||||
@@ -756,15 +756,15 @@ impl ECStore {
|
||||
|
||||
async fn pools_with_object(&self, pools: &[PoolObjInfo], opts: &ObjectOptions) -> Vec<PoolErr> {
|
||||
let mut errs = Vec::new();
|
||||
let pool_meta = self.pool_meta.read().await;
|
||||
|
||||
for pool in pools.iter() {
|
||||
if opts.skip_decommissioned && pool_meta.is_suspended(pool.index) {
|
||||
if opts.skip_decommissioned && self.is_suspended(pool.index).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
if opts.skip_rebalancing && self.is_pool_rebalancing(pool.index).await {
|
||||
continue;
|
||||
}
|
||||
// TODO:SkipRebalancing
|
||||
// if opts.SkipRebalancing && z.IsPoolRebalancing(pinfo.Index) {
|
||||
// continue
|
||||
// }
|
||||
|
||||
if let Some(err) = &pool.err {
|
||||
if is_err_read_quorum(err) {
|
||||
@@ -1865,7 +1865,9 @@ impl StorageAPI for ECStore {
|
||||
}
|
||||
|
||||
for pool in self.pools.iter() {
|
||||
// TODO: IsSuspended
|
||||
if self.is_suspended(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
|
||||
Ok(res) => return Ok(res),
|
||||
Err(err) => {
|
||||
@@ -1915,6 +1917,9 @@ impl StorageAPI for ECStore {
|
||||
let mut uploads = Vec::new();
|
||||
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
let res = pool
|
||||
.list_multipart_uploads(
|
||||
bucket,
|
||||
@@ -1948,7 +1953,9 @@ impl StorageAPI for ECStore {
|
||||
}
|
||||
|
||||
for (idx, pool) in self.pools.iter().enumerate() {
|
||||
// // TODO: IsSuspended
|
||||
if self.is_suspended(idx).await || self.is_pool_rebalancing(idx).await {
|
||||
continue;
|
||||
}
|
||||
let res = pool
|
||||
.list_multipart_uploads(bucket, object, None, None, None, MAX_UPLOADS_LIST)
|
||||
.await?;
|
||||
@@ -1983,8 +1990,8 @@ impl StorageAPI for ECStore {
|
||||
return self.pools[0].get_multipart_info(bucket, object, upload_id, opts).await;
|
||||
}
|
||||
|
||||
for (idx, pool) in self.pools.iter().enumerate() {
|
||||
if self.is_suspended(idx).await {
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -2017,7 +2024,9 @@ impl StorageAPI for ECStore {
|
||||
}
|
||||
|
||||
for pool in self.pools.iter() {
|
||||
// TODO: IsSuspended
|
||||
if self.is_suspended(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
let err = match pool.abort_multipart_upload(bucket, object, upload_id, opts).await {
|
||||
Ok(_) => return Ok(()),
|
||||
@@ -2060,7 +2069,9 @@ impl StorageAPI for ECStore {
|
||||
}
|
||||
|
||||
for pool in self.pools.iter() {
|
||||
// TODO: IsSuspended
|
||||
if self.is_suspended(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
let err = match pool
|
||||
.complete_multipart_upload(bucket, object, upload_id, uploaded_parts.clone(), opts)
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user