From 46c50b1bb1d99cf253c27de854f21199f1e2f2db Mon Sep 17 00:00:00 2001 From: likewu Date: Mon, 23 Jun 2025 18:51:31 +0800 Subject: [PATCH] Fix/main (#505) * fix error * fix * fix cargo --- Cargo.toml | 5 +- ecstore/Cargo.toml | 3 +- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 2 +- .../src/bucket/object_lock/objectlock_sys.rs | 4 +- ecstore/src/client/object_api_utils.rs | 4 +- ecstore/src/set_disk.rs | 50 +++++++++---------- reader/src/reader.rs | 6 ++- rustfs/src/admin/handlers/tier.rs | 10 ++-- rustfs/src/storage/ecfs.rs | 6 +-- 9 files changed, 44 insertions(+), 46 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index cccc3ebda..095dca8ec 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -22,6 +22,7 @@ members = [ "s3select/api", # S3 Select API interface "s3select/query", # S3 Select query engine "reader", + ] resolver = "2" @@ -90,7 +91,7 @@ derive_builder = "0.20.2" dioxus = { version = "0.6.3", features = ["router"] } dirs = "6.0.0" flatbuffers = "25.2.10" -flexi_logger = { version = "0.30.2", features = ["trc", "dont_minimize_extra_stacks"] } +flexi_logger = { version = "0.30.2", features = ["trc","dont_minimize_extra_stacks"] } form_urlencoded = "1.2.1" futures = "0.3.31" futures-core = "0.3.31" @@ -99,7 +100,6 @@ glob = "0.3.2" hex = "0.4.3" hex-simd = "0.8.0" highway = { version = "1.3.0" } -hmac = "0.12.1" hyper = "1.6.0" hyper-util = { version = "0.1.14", features = [ "tokio", @@ -201,6 +201,7 @@ serde-xml-rs = "0.8.1" serde_urlencoded = "0.7.1" sha1 = "0.10.6" sha2 = "0.10.9" +hmac = "0.12.1" std-next = "0.1.8" siphasher = "1.0.1" smallvec = { version = "1.15.1", features = ["serde"] } diff --git a/ecstore/Cargo.toml b/ecstore/Cargo.toml index 43bd4733a..dfb42f3bd 100644 --- a/ecstore/Cargo.toml +++ b/ecstore/Cargo.toml @@ -61,6 +61,7 @@ base64 = { workspace = true } hmac = { workspace = true } sha2 = { workspace = true } sha1 = { workspace = true } + hex-simd = { workspace = true } path-clean = { workspace = true } tempfile.workspace = true @@ -91,7 +92,7 @@ urlencoding = { workspace = true } smallvec = { workspace = true } shadow-rs.workspace = true rustfs-filemeta.workspace = true -rustfs-utils = { workspace = true, features = ["full"] } +rustfs-utils ={workspace = true, features=["full"]} rustfs-rio.workspace = true futures-util.workspace = true reader = { workspace = true } diff --git a/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 32942e1aa..501f50e56 100644 --- a/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -804,7 +804,7 @@ impl LifecycleOps for ObjectInfo { user_tags: self.user_tags.clone(), version_id: self.version_id.expect("err").to_string(), mod_time: self.mod_time, - size: self.size, + size: self.size as usize, is_latest: self.is_latest, num_versions: self.num_versions, delete_marker: self.delete_marker, diff --git a/ecstore/src/bucket/object_lock/objectlock_sys.rs b/ecstore/src/bucket/object_lock/objectlock_sys.rs index 5c71a093b..299b96989 100644 --- a/ecstore/src/bucket/object_lock/objectlock_sys.rs +++ b/ecstore/src/bucket/object_lock/objectlock_sys.rs @@ -36,7 +36,7 @@ pub fn enforce_retention_for_deletion(obj_info: &ObjectInfo) -> bool { return false; } - let lhold = objectlock::get_object_legalhold_meta(obj_info.user_defined.clone().expect("err")); + let lhold = objectlock::get_object_legalhold_meta(obj_info.user_defined.clone()); match lhold.status { Some(st) if st.as_str() == ObjectLockLegalHoldStatus::ON => { return true; @@ -44,7 +44,7 @@ pub fn enforce_retention_for_deletion(obj_info: &ObjectInfo) -> bool { _ => (), } - let ret = objectlock::get_object_retention_meta(obj_info.user_defined.clone().expect("err")); + let ret = objectlock::get_object_retention_meta(obj_info.user_defined.clone()); match ret.mode { Some(r) if (r.as_str() == ObjectLockRetentionMode::COMPLIANCE || r.as_str() == ObjectLockRetentionMode::GOVERNANCE) => { let t = objectlock::utc_now_ntp(); diff --git a/ecstore/src/client/object_api_utils.rs b/ecstore/src/client/object_api_utils.rs index e824968b8..911307b04 100644 --- a/ecstore/src/client/object_api_utils.rs +++ b/ecstore/src/client/object_api_utils.rs @@ -55,8 +55,8 @@ fn part_number_to_rangespec(oi: ObjectInfo, part_number: usize) -> Option { - if !cache.info.last_update.eq(&last_save) { - let _ = cache.save(DATA_USAGE_CACHE_NAME).await; - let _ = updates.send(cache.clone()).await; - } - } - result = buckets_results_rx.recv() => { - match result { - Some(result) => { - cache.replace(&result.name, &result.parent, result.entry); - cache.info.last_update = Some(SystemTime::now()); - }, - None => { - need_loop = false; - cache.info.next_cycle = want_cycle; - cache.info.last_update = Some(SystemTime::now()); + let last_save = Some(SystemTime::now()); + let mut need_loop = true; + while need_loop { + select! { + _ = ticker.tick() => { + if !cache.info.last_update.eq(&last_save) { let _ = cache.save(DATA_USAGE_CACHE_NAME).await; let _ = updates.send(cache.clone()).await; } } + result = buckets_results_rx.recv() => { + match result { + Some(result) => { + cache.replace(&result.name, &result.parent, result.entry); + cache.info.last_update = Some(SystemTime::now()); + }, + None => { + need_loop = false; + cache.info.next_cycle = want_cycle; + cache.info.last_update = Some(SystemTime::now()); + let _ = cache.save(DATA_USAGE_CACHE_NAME).await; + let _ = updates.send(cache.clone()).await; + } + } + } } } - } })) } else { None @@ -3183,7 +3183,7 @@ impl SetDisks { info!("ns_scanner start"); let _ = join_all(futures).await; if let Some(task) = task { - let _ = task.await; + let _ = task.await; } info!("ns_scanner completed"); Ok(()) @@ -3721,9 +3721,7 @@ impl SetDisks { let mut oi = obj_info.clone(); oi.metadata_only = true; - if let Some(user_defined) = &mut oi.user_defined { - user_defined.remove(X_AMZ_RESTORE.as_str()); - } + oi.user_defined.remove(X_AMZ_RESTORE.as_str()); let version_id = oi.version_id.clone().map(|v| v.to_string()); let obj = self @@ -4044,7 +4042,7 @@ impl ObjectIO for SetDisks { if opts.data_movement { fi.set_data_moved(); - } + } } let (online_disks, _, op_old_dir) = Self::rename_data( diff --git a/reader/src/reader.rs b/reader/src/reader.rs index e68a64aed..92c5ad14a 100644 --- a/reader/src/reader.rs +++ b/reader/src/reader.rs @@ -1,3 +1,5 @@ +#![allow(unused_imports)] + use bytes::Bytes; use s3s::StdError; use std::any::Any; @@ -180,7 +182,7 @@ pub trait Reader { async fn seek(&mut self, offset: usize) -> Result<()>; async fn read_exact(&mut self, buf: &mut [u8]) -> Result; async fn read_all(&mut self) -> Result> { - let mut data = Vec::new(); + let data = Vec::new(); Ok(data) } @@ -219,7 +221,7 @@ impl Reader for BufferReader { } #[tracing::instrument(level = "debug", skip(self))] async fn read_exact(&mut self, buf: &mut [u8]) -> Result { - let bytes_read = self.inner.read_exact(buf)?; + let _bytes_read = self.inner.read_exact(buf)?; self.pos += buf.len(); //Ok(bytes_read) Ok(0) diff --git a/rustfs/src/admin/handlers/tier.rs b/rustfs/src/admin/handlers/tier.rs index a11a92e71..2069268f5 100644 --- a/rustfs/src/admin/handlers/tier.rs +++ b/rustfs/src/admin/handlers/tier.rs @@ -1,7 +1,7 @@ use std::str::from_utf8; use http::{HeaderMap, StatusCode}; -use iam::get_global_action_cred; +//use iam::get_global_action_cred; use matchit::Params; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error}; use serde::Deserialize; @@ -277,8 +277,8 @@ impl Operation for RemoveTier { let (cred, _owner) = check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; - let sys_cred = get_global_action_cred() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "get_global_action_cred failed"))?; + //let sys_cred = get_global_action_cred() + // .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "get_global_action_cred failed"))?; let mut force: bool = false; let force_str = query.force; @@ -342,8 +342,8 @@ impl Operation for VerifyTier { let (cred, _owner) = check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; - let sys_cred = get_global_action_cred() - .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "get_global_action_cred failed"))?; + //let sys_cred = get_global_action_cred() + // .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "get_global_action_cred failed"))?; let mut tier_config_mgr = GLOBAL_TierConfigMgr.write().await; tier_config_mgr.verify(&query.tier.unwrap()); diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 98565da21..777e02d74 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -56,10 +56,6 @@ use ecstore::store_api::PutObjReader; use ecstore::store_api::StorageAPI; // use ecstore::store_api::RESERVED_METADATA_PREFIX; use ecstore::bucket::lifecycle::bucket_lifecycle_ops::validate_transition_tier; -use ecstore::bucket::utils::serialize; -use ecstore::cmd::bucket_replication::ReplicationStatusType; -use ecstore::cmd::bucket_replication::ReplicationType; -use ecstore::store_api::RESERVED_METADATA_PREFIX_LOWER; use futures::pin_mut; use futures::{Stream, StreamExt}; use http::HeaderMap; @@ -2458,7 +2454,7 @@ impl S3 for FS { let retain_until_date = object_info .user_defined .get("x-amz-object-lock-retain-until-date") - .and_then(|v| OffsetDateTime::parse(v.as_str(), &Rfc3339).ok()) + .and_then(|v| OffsetDateTime::parse(v.as_str(), &Rfc3339).ok()) .map(Timestamp::from); Ok(S3Response::new(GetObjectRetentionOutput {