From 750194e4dc2e45c4d7f1eb585fd1490cad9aa7dc Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 28 Apr 2025 13:16:11 +0800 Subject: [PATCH 1/6] test --- Cargo.toml | 2 +- ecstore/src/cache_value/metacache_set.rs | 15 ++- ecstore/src/disk/mod.rs | 147 +++++++++++++++-------- ecstore/src/file_meta.rs | 13 +- ecstore/src/heal/data_scanner.rs | 2 +- ecstore/src/pools.rs | 9 +- ecstore/src/rebalance.rs | 8 +- ecstore/src/set_disk.rs | 4 +- ecstore/src/store_list_objects.rs | 4 +- scripts/run.sh | 2 +- 10 files changed, 130 insertions(+), 76 deletions(-) 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" From 26d69cdc7f31b5ca8ba63e9f388a27d2572b0f70 Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 28 Apr 2025 16:58:37 +0800 Subject: [PATCH 2/6] test --- ecstore/src/cache_value/metacache_set.rs | 10 +---- ecstore/src/pools.rs | 34 +++++++++++----- ecstore/src/rebalance.rs | 14 +++---- ecstore/src/store.rs | 49 +++++++++++++++--------- 4 files changed, 61 insertions(+), 46 deletions(-) diff --git a/ecstore/src/cache_value/metacache_set.rs b/ecstore/src/cache_value/metacache_set.rs index c3e75d093..1e32b4391 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, warn}; +use tracing::error; pub type AgreedFn = Box Pin + Send>> + Send + 'static>; pub type PartialFn = Box]) -> Pin + Send>> + Send + 'static>; @@ -54,7 +54,6 @@ impl Clone for ListPathRawOptions { pub async fn list_path_raw(mut rx: B_Receiver, 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")); } @@ -214,16 +213,12 @@ pub async fn list_path_raw(mut rx: B_Receiver, 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()]; for item in top_entries.iter_mut().take(i) { *item = None; @@ -270,7 +265,6 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - } } - warn!("list_path_raw: all at eof or error"); break; } @@ -279,7 +273,6 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - 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; } @@ -293,7 +286,6 @@ 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/pools.rs b/ecstore/src/pools.rs index c151b0627..8a0761006 100644 --- a/ecstore/src/pools.rs +++ b/ecstore/src/pools.rs @@ -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, 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,6 +1340,7 @@ 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]| { let resolver = resolver.clone(); let cb_func = cb_func.clone(); diff --git a/ecstore/src/rebalance.rs b/ecstore/src/rebalance.rs index d6ec57113..97255caf8 100644 --- a/ecstore/src/rebalance.rs +++ b/ecstore/src/rebalance.rs @@ -370,7 +370,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); } @@ -581,17 +581,17 @@ 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: {}", 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; } @@ -599,7 +599,7 @@ impl ECStore { 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(); @@ -838,8 +838,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()) @@ -1068,7 +1066,7 @@ 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]| { // let cb = cb.clone(); let resolver = resolver.clone(); diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 5a2075775..821eb5305 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -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 { 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) From 738a81e50e992419588f3d473cde8b63119fb8fe Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 28 Apr 2025 17:35:27 +0800 Subject: [PATCH 3/6] test --- ecstore/src/rebalance.rs | 29 ++++++++++++++++++----------- 1 file changed, 18 insertions(+), 11 deletions(-) diff --git a/ecstore/src/rebalance.rs b/ecstore/src/rebalance.rs index 97255caf8..007cf7a4b 100644 --- a/ecstore/src/rebalance.rs +++ b/ecstore/src/rebalance.rs @@ -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}; @@ -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); } } } @@ -584,7 +590,7 @@ impl ECStore { warn!("Pool {} rebalancing is started", pool_index + 1); while let Some(bucket) = self.next_rebal_bucket(pool_index).await? { - warn!("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") { @@ -596,6 +602,7 @@ impl ECStore { break; } + warn!("Rebalance bucket: done {} ", bucket); self.bucket_rebalance_done(pool_index, bucket).await?; } @@ -632,7 +639,6 @@ impl ECStore { false } - #[allow(unused_assignments)] #[tracing::instrument(skip(self, wk, set))] async fn rebalance_entry( &self, @@ -854,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( @@ -881,9 +893,6 @@ impl ECStore { part_num: pi.part_num, e_tag: pi.etag, }; - - // 把reader所有权拿回来? - reader = data.stream; } if let Err(err) = self @@ -937,7 +946,7 @@ impl ECStore { #[tracing::instrument(skip(self, rx))] async fn rebalance_bucket(self: &Arc, rx: B_Receiver, 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{ @@ -975,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; }); } From cdcf5d091759e1535e7ea97fe580194945ca9a47 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Mon, 28 Apr 2025 06:25:20 +0000 Subject: [PATCH 4/6] fix readme Signed-off-by: junxiang Mu <1948535941@qq.com> --- README.md | 2 +- ecstore/src/erasure.rs | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/README.md b/README.md index 511f6736e..812f13a3c 100644 --- a/README.md +++ b/README.md @@ -84,7 +84,7 @@ observability data formats (e.g. Jaeger, Prometheus, etc.) sending to one or mor 2. Run the following command: ```bash -docker-compose -f docker-compose.yml up -d +docker compose -f docker-compose.yml up -d ``` 3. Access the Grafana dashboard by navigating to `http://localhost:3000` in your browser. The default username and diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index 021ae7bed..590889e78 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -470,7 +470,7 @@ impl Erasure { self.encoder.as_ref().unwrap().reconstruct(&mut bufs)?; } - let shards = bufs.into_iter().flatten().collect::>(); + let shards = bufs.into_iter().flatten().map(Bytes::from).collect::>(); if shards.len() != self.parity_shards + self.data_shards { return Err(Error::from_string("can not reconstruct data")); } @@ -479,7 +479,7 @@ impl Erasure { if w.is_none() { continue; } - match w.as_mut().unwrap().write(shards[i].clone().into()).await { + match w.as_mut().unwrap().write(shards[i].clone()).await { Ok(_) => {} Err(e) => { info!("write failed, err: {:?}", e); From 6aecd72acc6219ae791d9a069853b88c075768c6 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 28 Apr 2025 14:37:28 +0800 Subject: [PATCH 5/6] improve readme.md --- .docker/observability/README.md | 16 ++++++++++++---- deploy/README.md | 18 +++++++++++++----- deploy/build/rustfs-zh.service | 4 ++++ deploy/build/rustfs.service | 4 ++++ deploy/certs/README.md | 16 +++++++++++++--- .../config/{.example.env => .example.obs.env} | 0 6 files changed, 46 insertions(+), 12 deletions(-) rename deploy/config/{.example.env => .example.obs.env} (100%) diff --git a/.docker/observability/README.md b/.docker/observability/README.md index a84f7592d..3d40319b1 100644 --- a/.docker/observability/README.md +++ b/.docker/observability/README.md @@ -6,7 +6,7 @@ This directory contains the observability stack for the application. The stack i - Grafana 11.6.0 - Loki 3.4.2 - Jaeger 2.4.0 -- Otel Collector 0.120.0 #0.121.0 remove loki +- Otel Collector 0.120.0 # 0.121.0 remove loki ## Prometheus @@ -47,8 +47,16 @@ observability data formats (e.g. Jaeger, Prometheus, etc.) sending to one or mor To deploy the observability stack, run the following command: +- docker latest version + ```bash -docker-compose -f docker-compose.yml -f docker-compose.override.yml up -d +docker compose -f docker-compose.yml -f docker-compose.override.yml up -d +``` + +- docker compose v2.0.0 or before + +```bash +docke-compose -f docker-compose.yml -f docker-compose.override.yml up -d ``` To access the Grafana dashboard, navigate to `http://localhost:3000` in your browser. The default username and password @@ -63,7 +71,7 @@ To access the Prometheus dashboard, navigate to `http://localhost:9090` in your To stop the observability stack, run the following command: ```bash -docker-compose -f docker-compose.yml -f docker-compose.override.yml down +docker compose -f docker-compose.yml -f docker-compose.override.yml down ``` ## How to remove data @@ -71,7 +79,7 @@ docker-compose -f docker-compose.yml -f docker-compose.override.yml down To remove the data generated by the observability stack, run the following command: ```bash -docker-compose -f docker-compose.yml -f docker-compose.override.yml down -v +docker compose -f docker-compose.yml -f docker-compose.override.yml down -v ``` ## How to configure diff --git a/deploy/README.md b/deploy/README.md index 4d3d4476a..2efdd85ec 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -23,13 +23,21 @@ managing and monitoring the system. | |--rustfs.service // systemd service file | |--rustfs-zh.service.md // systemd service file in Chinese |--certs -| |--README.md // certs readme -| |--rustfs_tls_cert.pem // API cert.pem -| |--rustfs_tls_key.pem // API key.pem -| |--rustfs_console_tls_cert.pem // console cert.pem -| |--rustfs_console_tls_key.pem // console key.pem +| ├── rustfs_cert.pem // Default|fallback certificate +| ├── rustfs_key.pem // Default|fallback private key +| ├── example.com/ // certificate directory of specific domain names +| │ ├── rustfs_cert.pem +| │ └── rustfs_key.pem +| ├── api.example.com/ +| │ ├── rustfs_cert.pem +| │ └── rustfs_key.pem +| └── cdn.example.com/ +| ├── rustfs_cert.pem +| └── rustfs_key.pem |--config | |--obs.example.yaml // example config | |--rustfs.env // env config | |--rustfs-zh.env // env config in Chinese +| |--.example.obs.env // example env config +| |--event.example.toml // event config ``` \ No newline at end of file diff --git a/deploy/build/rustfs-zh.service b/deploy/build/rustfs-zh.service index a1cd0df05..17351e675 100644 --- a/deploy/build/rustfs-zh.service +++ b/deploy/build/rustfs-zh.service @@ -51,6 +51,10 @@ ExecStart=/usr/local/bin/rustfs \ EnvironmentFile=-/etc/default/rustfs ExecStart=/usr/local/bin/rustfs $RUSTFS_VOLUMES $RUSTFS_OPTS +# standard output and error log configuration +StandardOutput=append:/data/deploy/rust/logs/rustfs.log +StandardError=append:/data/deploy/rust/logs/rustfs-err.log + # resource constraints LimitNOFILE=1048576 # 设置文件描述符上限为 1048576,支持高并发连接。 diff --git a/deploy/build/rustfs.service b/deploy/build/rustfs.service index df6e4067b..9c72e4276 100644 --- a/deploy/build/rustfs.service +++ b/deploy/build/rustfs.service @@ -31,6 +31,10 @@ ExecStart=/usr/local/bin/rustfs \ EnvironmentFile=-/etc/default/rustfs ExecStart=/usr/local/bin/rustfs $RUSTFS_VOLUMES $RUSTFS_OPTS +# service log configuration +StandardOutput=append:/data/deploy/rust/logs/rustfs.log +StandardError=append:/data/deploy/rust/logs/rustfs-err.log + # resource constraints LimitNOFILE=1048576 LimitNPROC=32768 diff --git a/deploy/certs/README.md b/deploy/certs/README.md index 84b733e79..e36d188b4 100644 --- a/deploy/certs/README.md +++ b/deploy/certs/README.md @@ -32,7 +32,17 @@ openssl req -x509 -newkey rsa:2048 -keyout key.pem -out cert.pem -days 365 -node ### TLS File ```text - rustfs_public.crt #api cert.pem - - rustfs_private.key #api key.pem +cd deploy/certs/ +ls -la + ├── rustfs_cert.pem // Default|fallback certificate + ├── rustfs_key.pem // Default|fallback private key + ├── example.com/ // certificate directory of specific domain names + │ ├── rustfs_cert.pem + │ └── rustfs_key.pem + ├── api.example.com/ + │ ├── rustfs_cert.pem + │ └── rustfs_key.pem + └── cdn.example.com/ + ├── rustfs_cert.pem + └── rustfs_key.pem ``` \ No newline at end of file diff --git a/deploy/config/.example.env b/deploy/config/.example.obs.env similarity index 100% rename from deploy/config/.example.env rename to deploy/config/.example.obs.env From 79d58a98f413498ec01bb132b32f4dc3f1a5c157 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 28 Apr 2025 14:51:51 +0800 Subject: [PATCH 6/6] improve code for readme.md add chinese readme.md --- .docker/observability/README_ZH.md | 42 +++++++++ README.md | 141 ++++++++++++++++------------- README_ZH.md | 136 ++++++++++++++++++++++++++++ 3 files changed, 255 insertions(+), 64 deletions(-) create mode 100644 .docker/observability/README_ZH.md create mode 100644 README_ZH.md diff --git a/.docker/observability/README_ZH.md b/.docker/observability/README_ZH.md new file mode 100644 index 000000000..7ba5342bf --- /dev/null +++ b/.docker/observability/README_ZH.md @@ -0,0 +1,42 @@ +## 部署可观测性系统 + +OpenTelemetry Collector 提供了一个厂商中立的遥测数据处理方案,用于接收、处理和导出遥测数据。它消除了为支持多种开源可观测性数据格式(如 +Jaeger、Prometheus 等)而需要运行和维护多个代理/收集器的必要性。 + +### 快速部署 + +1. 进入 `.docker/observability` 目录 +2. 执行以下命令启动服务: + +```bash +docker compose up -d -f docker-compose.yml +``` + +### 访问监控面板 + +服务启动后,可通过以下地址访问各个监控面板: + +- Grafana: `http://localhost:3000` (默认账号/密码:`admin`/`admin`) +- Jaeger: `http://localhost:16686` +- Prometheus: `http://localhost:9090` + +## 配置可观测性 + +### 创建配置文件 + +1. 进入 `deploy/config` 目录 +2. 复制示例配置:`cp obs.toml.example obs.toml` +3. 编辑 `obs.toml` 配置文件,修改以下关键参数: + +| 配置项 | 说明 | 示例值 | +|-----------------|----------------------------|-----------------------| +| endpoint | OpenTelemetry Collector 地址 | http://localhost:4317 | +| service_name | 服务名称 | rustfs | +| service_version | 服务版本 | 1.0.0 | +| environment | 运行环境 | production | +| meter_interval | 指标导出间隔 (秒) | 30 | +| sample_ratio | 采样率 | 1.0 | +| use_stdout | 是否输出到控制台 | true/false | +| logger_level | 日志级别 | info | + +``` \ No newline at end of file diff --git a/README.md b/README.md index 812f13a3c..79d9a342c 100644 --- a/README.md +++ b/README.md @@ -1,54 +1,68 @@ -# How to compile RustFS +# RustFS -| Must package | Version | download link | -|--------------|---------|----------------------------------------------------------------------------------------------------------------------------------| -| Rust | 1.8.5 | https://www.rust-lang.org/tools/install | -| protoc | 30.2 | [protoc-30.2-linux-x86_64.zip](https://github.com/protocolbuffers/protobuf/releases/download/v30.2/protoc-30.2-linux-x86_64.zip) | -| flatc | 24.0+ | [Linux.flatc.binary.g++-13.zip](https://github.com/google/flatbuffers/releases/download/v25.2.10/Linux.flatc.binary.g++-13.zip) | +## English Documentation |[中文文档](README_ZH.md) -Download Links: +### Prerequisites -https://github.com/google/flatbuffers/releases/download/v25.2.10/Linux.flatc.binary.g++-13.zip +| Package | Version | Download Link | +|---------|---------|----------------------------------------------------------------------------------------------------------------------------------| +| Rust | 1.8.5+ | [rust-lang.org/tools/install](https://www.rust-lang.org/tools/install) | +| protoc | 30.2+ | [protoc-30.2-linux-x86_64.zip](https://github.com/protocolbuffers/protobuf/releases/download/v30.2/protoc-30.2-linux-x86_64.zip) | +| flatc | 24.0+ | [Linux.flatc.binary.g++-13.zip](https://github.com/google/flatbuffers/releases/download/v25.2.10/Linux.flatc.binary.g++-13.zip) | -https://github.com/protocolbuffers/protobuf/releases/download/v30.2/protoc-30.2-linux-x86_64.zip +### Building RustFS -generate protobuf code: +#### Generate Protobuf Code -```cargo run --bin gproto``` +```bash +cargo run --bin gproto +``` -Or use Docker: +#### Using Docker for Prerequisites -```yml +```yaml - uses: arduino/setup-protoc@v3 with: - version: "30.2" + version: "30.2" - uses: Nugine/setup-flatc@v1 with: - version: "25.2.10" + version: "25.2.10" ``` -# How to add Console web +#### Adding Console Web UI -1. `wget https://dl.rustfs.com/artifacts/console/rustfs-console-latest.zip` +1. Download the latest console UI: + ```bash + wget https://dl.rustfs.com/artifacts/console/rustfs-console-latest.zip + ``` +2. Create the static directory: + ```bash + mkdir -p ./rustfs/static + ``` +3. Extract and compile RustFS: + ```bash + unzip rustfs-console-latest.zip -d ./rustfs/static + cargo build + ``` -2. mkdir in this repos folder `./rustfs/static` +### Running RustFS -3. Compile RustFS +#### Configuration -# Star RustFS +Set the required environment variables: -Add Env Information: - -``` +```bash +# Basic config export RUSTFS_VOLUMES="./target/volume/test" export RUSTFS_ADDRESS="0.0.0.0:9000" export RUSTFS_CONSOLE_ENABLE=true export RUSTFS_CONSOLE_ADDRESS="0.0.0.0:9001" -# 具体路径修改为配置文件真实路径,obs.example.toml 仅供参考 其中`RUSTFS_OBS_CONFIG` 和下面变量二选一 -export RUSTFS_OBS_CONFIG="./deploy/config/obs.example.toml" -# 如下变量需要必须参数都有值才可以,以及会覆盖配置文件`obs.example.toml`中的值 +# Observability config (option 1: config file) +export RUSTFS_OBS_CONFIG="./deploy/config/obs.toml" + +# Observability config (option 2: environment variables) export RUSTFS__OBSERVABILITY__ENDPOINT=http://localhost:4317 export RUSTFS__OBSERVABILITY__USE_STDOUT=true export RUSTFS__OBSERVABILITY__SAMPLE_RATIO=2.0 @@ -57,6 +71,9 @@ export RUSTFS__OBSERVABILITY__SERVICE_NAME=rustfs export RUSTFS__OBSERVABILITY__SERVICE_VERSION=0.1.0 export RUSTFS__OBSERVABILITY__ENVIRONMENT=develop export RUSTFS__OBSERVABILITY__LOGGER_LEVEL=info +export RUSTFS__OBSERVABILITY__LOCAL_LOGGING_ENABLED=true + +# Logging sinks export RUSTFS__SINKS__FILE__ENABLED=true export RUSTFS__SINKS__FILE__PATH="./deploy/logs/rustfs.log" export RUSTFS__SINKS__WEBHOOK__ENABLED=false @@ -68,54 +85,50 @@ export RUSTFS__SINKS__KAFKA__TOPIC="" export RUSTFS__LOGGER__QUEUE_CAPACITY=10 ``` -You need replace your real data folder: +#### Start the service -``` +```bash ./rustfs /data/rustfs ``` -## How to deploy the observability stack +### Observability Stack -The OpenTelemetry Collector offers a vendor-agnostic implementation on how to receive, process, and export telemetry -data. It removes the need to run, operate, and maintain multiple agents/collectors in order to support open-source -observability data formats (e.g. Jaeger, Prometheus, etc.) sending to one or more open-source or commercial back-ends. +#### Deployment -1. Enter the `.docker/observability` directory, -2. Run the following command: +1. Navigate to the observability directory: + ```bash + cd .docker/observability + ``` -```bash -docker compose -f docker-compose.yml up -d -``` +2. Start the observability stack: + ```bash + docker compose up -d -f docker-compose.yml + ``` -3. Access the Grafana dashboard by navigating to `http://localhost:3000` in your browser. The default username and - password are `admin` and `admin`, respectively. +#### Access Monitoring Dashboards -4. Access the Jaeger dashboard by navigating to `http://localhost:16686` in your browser. +- Grafana: `http://localhost:3000` (credentials: `admin`/`admin`) +- Jaeger: `http://localhost:16686` +- Prometheus: `http://localhost:9090` -5. Access the Prometheus dashboard by navigating to `http://localhost:9090` in your browser. +#### Configuring Observability -## Create a new Observability configuration file - -#### 1. Enter the `deploy/config` directory, - -#### 2. Copy `obs.toml.example` to `obs.toml` - -#### 3. Modify the `obs.toml` configuration file - -##### 3.1. Modify the `endpoint` value to the address of the OpenTelemetry Collector - -##### 3.2. Modify the `service_name` value to the name of the service - -##### 3.3. Modify the `service_version` value to the version of the service - -##### 3.4. Modify the `environment` value to the environment of the service - -##### 3.5. Modify the `meter_interval` value to export interval - -##### 3.6. Modify the `sample_ratio` value to the sample ratio - -##### 3.7. Modify the `use_stdout` value to export to stdout - -##### 3.8. Modify the `logger_level` value to the logger level +1. Copy the example configuration: + ```bash + cd deploy/config + cp obs.toml.example obs.toml + ``` +2. Edit `obs.toml` with the following parameters: +| Parameter | Description | Example | +|----------------------|-----------------------------------|-----------------------| +| endpoint | OpenTelemetry Collector address | http://localhost:4317 | +| service_name | Service name | rustfs | +| service_version | Service version | 1.0.0 | +| environment | Runtime environment | production | +| meter_interval | Metrics export interval (seconds) | 30 | +| sample_ratio | Sampling ratio | 1.0 | +| use_stdout | Output to console | true/false | +| logger_level | Log level | info | +| local_logging_enable | stdout | true/false | diff --git a/README_ZH.md b/README_ZH.md new file mode 100644 index 000000000..c2016fb7c --- /dev/null +++ b/README_ZH.md @@ -0,0 +1,136 @@ +# RustFS + +## [English Documentation](README.md) |中文文档 + +### 前置要求 + +| 软件包 | 版本 | 下载链接 | +|--------|--------|----------------------------------------------------------------------------------------------------------------------------------| +| Rust | 1.8.5+ | [rust-lang.org/tools/install](https://www.rust-lang.org/tools/install) | +| protoc | 30.2+ | [protoc-30.2-linux-x86_64.zip](https://github.com/protocolbuffers/protobuf/releases/download/v30.2/protoc-30.2-linux-x86_64.zip) | +| flatc | 24.0+ | [Linux.flatc.binary.g++-13.zip](https://github.com/google/flatbuffers/releases/download/v25.2.10/Linux.flatc.binary.g++-13.zip) | + +### 构建 RustFS + +#### 生成 Protobuf 代码 + +```bash +cargo run --bin gproto +``` + +#### 使用 Docker 安装依赖 + +```yaml +- uses: arduino/setup-protoc@v3 + with: + version: "30.2" + +- uses: Nugine/setup-flatc@v1 + with: + version: "25.2.10" +``` + +#### 添加控制台 Web UI + +1. 下载最新的控制台 UI: + ```bash + wget https://dl.rustfs.com/artifacts/console/rustfs-console-latest.zip + ``` +2. 创建静态资源目录: + ```bash + mkdir -p ./rustfs/static + ``` +3. 解压并编译 RustFS: + ```bash + unzip rustfs-console-latest.zip -d ./rustfs/static + cargo build + ``` + +### 运行 RustFS + +#### 配置 + +设置必要的环境变量: + +```bash +# 基础配置 +export RUSTFS_VOLUMES="./target/volume/test" +export RUSTFS_ADDRESS="0.0.0.0:9000" +export RUSTFS_CONSOLE_ENABLE=true +export RUSTFS_CONSOLE_ADDRESS="0.0.0.0:9001" + +# 可观测性配置(方式一:配置文件) +export RUSTFS_OBS_CONFIG="./deploy/config/obs.toml" + +# 可观测性配置(方式二:环境变量) +export RUSTFS__OBSERVABILITY__ENDPOINT=http://localhost:4317 +export RUSTFS__OBSERVABILITY__USE_STDOUT=true +export RUSTFS__OBSERVABILITY__SAMPLE_RATIO=2.0 +export RUSTFS__OBSERVABILITY__METER_INTERVAL=30 +export RUSTFS__OBSERVABILITY__SERVICE_NAME=rustfs +export RUSTFS__OBSERVABILITY__SERVICE_VERSION=0.1.0 +export RUSTFS__OBSERVABILITY__ENVIRONMENT=develop +export RUSTFS__OBSERVABILITY__LOGGER_LEVEL=info +export RUSTFS__OBSERVABILITY__LOCAL_LOGGING_ENABLED=true + +# 日志接收器 +export RUSTFS__SINKS__FILE__ENABLED=true +export RUSTFS__SINKS__FILE__PATH="./deploy/logs/rustfs.log" +export RUSTFS__SINKS__WEBHOOK__ENABLED=false +export RUSTFS__SINKS__WEBHOOK__ENDPOINT="" +export RUSTFS__SINKS__WEBHOOK__AUTH_TOKEN="" +export RUSTFS__SINKS__KAFKA__ENABLED=false +export RUSTFS__SINKS__KAFKA__BOOTSTRAP_SERVERS="" +export RUSTFS__SINKS__KAFKA__TOPIC="" +export RUSTFS__LOGGER__QUEUE_CAPACITY=10 +``` + +#### 启动服务 + +```bash +./rustfs /data/rustfs +``` + +### 可观测性系统 + +#### 部署 + +1. 进入可观测性目录: + ```bash + cd .docker/observability + ``` + +2. 启动可观测性系统: + ```bash + docker compose up -d -f docker-compose.yml + ``` + +#### 访问监控面板 + +- Grafana: `http://localhost:3000` (默认账号/密码:`admin`/`admin`) +- Jaeger: `http://localhost:16686` +- Prometheus: `http://localhost:9090` + +#### 配置可观测性 + +1. 复制示例配置: + ```bash + cd deploy/config + cp obs.toml.example obs.toml + ``` + +2. 编辑 `obs.toml` 配置文件,参数如下: + +| 配置项 | 说明 | 示例值 | +|----------------------|----------------------------|-----------------------| +| endpoint | OpenTelemetry Collector 地址 | http://localhost:4317 | +| service_name | 服务名称 | rustfs | +| service_version | 服务版本 | 1.0.0 | +| environment | 运行环境 | production | +| meter_interval | 指标导出间隔 (秒) | 30 | +| sample_ratio | 采样率 | 1.0 | +| use_stdout | 是否输出到控制台 | true/false | +| logger_level | 日志级别 | info | +| local_logging_enable | 控制台是否答应日志 | true/false | + +``` \ No newline at end of file