diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 4d89a5320..f6e7df9d2 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -295,7 +295,6 @@ impl BucketMetadataSys { let mut futures = Vec::new(); for bucket in buckets.iter() { - // TODO: HealBucket let api = self.api.clone(); let bucket = bucket.clone(); futures.push(async move { @@ -319,14 +318,13 @@ impl BucketMetadataSys { let mut mp = self.metadata_map.write().await; - // TODO:EventNotifier,BucketTargetSys + // TODO: EventNotifier for res in results { match res { Ok(res) => { if let Some(bucket) = buckets.get(idx) { let x = Arc::new(res); mp.insert(bucket.clone(), x.clone()); - // TODO:EventNotifier,BucketTargetSys BucketTargetSys::get().set(bucket, &x).await; } } diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index ee75fac1b..9d8d26526 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -918,7 +918,6 @@ fn decode_batch_read_version_response_items( Ok(batch_read_version_resps) } -// TODO: all api need to handle errors #[async_trait::async_trait] impl DiskAPI for RemoteDisk { #[tracing::instrument(skip(self))] @@ -1185,56 +1184,6 @@ impl DiskAPI for RemoteDisk { .await } - // // FIXME: TODO: use writer - // #[tracing::instrument(skip(self, wr))] - // async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()> { - // let now = std::time::SystemTime::now(); - // info!("walk_dir {}/{}/{:?}", self.endpoint.to_string(), opts.bucket, opts.filter_prefix); - // let mut wr = wr; - // let mut out = MetacacheWriter::new(&mut wr); - // let mut buf = Vec::new(); - // opts.serialize(&mut Serializer::new(&mut buf))?; - // let mut client = node_service_time_out_client(&self.addr) - // .await - // .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?; - // let request = Request::new(WalkDirRequest { - // disk: self.endpoint.to_string(), - // walk_dir_options: buf.into(), - // }); - // let mut response = client.walk_dir(request).await?.into_inner(); - - // loop { - // match response.next().await { - // Some(Ok(resp)) => { - // if !resp.success { - // if let Some(err) = resp.error_info { - // if err == "Unexpected EOF" { - // return Err(Error::Io(std::io::Error::new(std::io::ErrorKind::UnexpectedEof, err))); - // } else { - // return Err(Error::other(err)); - // } - // } - - // return Err(Error::other("unknown error")); - // } - // let entry = serde_json::from_str::(&resp.meta_cache_entry) - // .map_err(|_| Error::other(format!("Unexpected response: {:?}", response)))?; - // out.write_obj(&entry).await?; - // } - // None => break, - // _ => return Err(Error::other(format!("Unexpected response: {:?}", response))), - // } - // } - - // info!( - // "walk_dir {}/{:?} done {:?}", - // opts.bucket, - // opts.filter_prefix, - // now.elapsed().unwrap_or_default() - // ); - // Ok(()) - // } - #[tracing::instrument(skip(self))] async fn delete_version( &self, diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 5f765621d..6d6ebda67 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -176,7 +176,6 @@ impl Sets { let sets = Arc::new(Self { id: fm.id, - // sets: todo!(), disk_set, pool_idx, endpoints: endpoints.clone(), @@ -194,10 +193,6 @@ impl Sets { let rx1 = rx.resubscribe(); tokio::spawn(async move { asets.monitor_and_connect_endpoints(rx1).await }); - // let sets2 = sets.clone(); - // let rx2 = rx.resubscribe(); - // tokio::spawn(async move { sets2.cleanup_deleted_objects_loop(rx2).await }); - Ok(sets) } @@ -205,36 +200,6 @@ impl Sets { self.set_drive_count } - // pub async fn cleanup_deleted_objects_loop(self: Arc, mut rx: Receiver<()>) { - // tokio::time::sleep(Duration::from_secs(5)).await; - - // info!("start cleanup_deleted_objects_loop"); - - // // TODO: config interval - // let mut interval = tokio::time::interval(Duration::from_secs(15 * 3)); - // loop { - // tokio::select! { - // _= interval.tick()=>{ - - // info!("cleanup_deleted_objects_loop tick"); - - // for set in self.disk_set.iter() { - // set.clone().cleanup_deleted_objects().await; - // } - - // interval.reset(); - // }, - - // _ = rx.recv() => { - // warn!("cleanup_deleted_objects_loop ctx cancelled"); - // break; - // } - // } - // } - - // warn!("cleanup_deleted_objects_loop exit"); - // } - pub async fn monitor_and_connect_endpoints(&self, mut rx: Receiver<()>) { tokio::time::sleep(Duration::from_secs(5)).await; @@ -327,17 +292,6 @@ impl Sets { } } - // async fn commit_rename_data_dir( - // &self, - // disks: &Vec>, - // bucket: &str, - // object: &str, - // data_dir: &str, - // // write_quorum: usize, - // ) -> Vec> { - // unimplemented!() - // } - async fn delete_prefix(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { let mut futures = Vec::new(); let mut opt = opts.clone(); diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index e876f7a35..7bf9649ec 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -1785,14 +1785,7 @@ impl LocalDisk { DiskMetrics::default() } - async fn bitrot_verify( - &self, - part_path: &PathBuf, - part_size: usize, - algo: HashAlgorithm, - sum: &[u8], - shard_size: usize, - ) -> Result<()> { + async fn bitrot_verify(&self, part_path: &PathBuf, part_size: usize, algo: HashAlgorithm, shard_size: usize) -> Result<()> { let retry_count = bitrot_size_mismatch_retry_count(); let retry_delay = bitrot_size_mismatch_retry_delay(); @@ -1801,16 +1794,7 @@ impl LocalDisk { let meta = file.metadata().await.map_err(to_file_error)?; let file_size = meta.len() as usize; - match bitrot_verify( - Box::new(file), - file_size, - part_size, - algo.clone(), - Bytes::copy_from_slice(sum), - shard_size, - ) - .await - { + match bitrot_verify(Box::new(file), file_size, part_size, algo.clone(), shard_size).await { Ok(()) => return Ok(()), Err(err) if attempt < retry_count && is_bitrot_size_mismatch_error(&err) => { info!( @@ -2622,7 +2606,6 @@ impl DiskAPI for LocalDisk { &part_path, erasure.shard_file_size(part.size as i64) as usize, checksum_algo, - &checksum_info.hash, erasure.shard_size(), ) .await diff --git a/crates/ecstore/src/erasure/coding/bitrot.rs b/crates/ecstore/src/erasure/coding/bitrot.rs index 874ef052f..729317614 100644 --- a/crates/ecstore/src/erasure/coding/bitrot.rs +++ b/crates/ecstore/src/erasure/coding/bitrot.rs @@ -12,7 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -use bytes::Bytes; use pin_project_lite::pin_project; use rustfs_utils::HashAlgorithm; use std::io::IoSlice; @@ -221,7 +220,6 @@ pub async fn bitrot_verify( want_size: usize, part_size: usize, algo: HashAlgorithm, - _want: Bytes, // FIXME: useless parameter? mut shard_size: usize, ) -> std::io::Result<()> { let mut hash_buf = vec![0; algo.size()]; diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 26e758544..4f4cca1d6 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -365,7 +365,6 @@ impl ObjectInfo { if let Some(size_str) = rustfs_utils::http::get_str(&self.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE) && !size_str.is_empty() { - // Todo: deal with error let size = size_str.parse::().map_err(|e| std::io::Error::other(e.to_string()))?; return Ok(size); } @@ -509,8 +508,6 @@ impl ObjectInfo { }) .collect::>(); - // TODO: part checksums - ObjectInfo { bucket: bucket.to_string(), name, diff --git a/crates/ecstore/src/services/event_notification.rs b/crates/ecstore/src/services/event_notification.rs index b374bd9e9..a5a5c7acd 100644 --- a/crates/ecstore/src/services/event_notification.rs +++ b/crates/ecstore/src/services/event_notification.rs @@ -52,11 +52,6 @@ impl EventNotifier { } fn init_bucket_targets(&self, _api: ECStore) -> Result<(), std::io::Error> { - /*if err := self.target_list.Add(globalNotifyTargetList.Targets()...); err != nil { - return err - } - self.target_list = self.target_list.Init(runtime.GOMAXPROCS(0)) // TODO: make this configurable (y4m4) - nil*/ warn!("init_bucket_targets called but currently no-op in this build"); Ok(()) } diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index ff048ab66..52bacd69f 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -356,11 +356,6 @@ impl ECStore { "Rebalance bucket entry scan started" ); - // TODO: other config - // if bucket != RUSTFS_META_BUCKET{ - - // } - let pool = clone_arc_by_index(self.pools.as_slice(), pool_index, "invalid rebalance pool index")?; let bucket_configs = Arc::new(load_rebalance_bucket_configs(&bucket).await?); diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 8940c26bc..ab977e19a 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2476,7 +2476,6 @@ impl SetDisks { let mut futures = Vec::with_capacity(disks.len()); if let Some(ret_err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - // TODO: add concurrency for (i, err) in errs.iter().enumerate() { if err.is_some() { continue; @@ -2562,19 +2561,6 @@ impl SetDisks { vec![None; disks.len()] }; - // // TODO: reduce_common_data_dir - // if let Some(old_dir) = rename_ress - // .iter() - // .filter_map(|v| if v.is_some() { v.as_ref().unwrap().old_data_dir } else { None }) - // .map(|v| v.to_string()) - // .next() - // { - // let cm_errs = self.commit_rename_data_dir(&shuffle_disks, &bucket, &object, &old_dir).await; - // warn!("put_object commit_rename_data_dir:{:?}", &cm_errs); - // } - - // self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_dir).await?; - Ok((online_disks, versions, data_dir, cleanup_disks)) } diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 650e1affb..ae6823b50 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -1901,14 +1901,6 @@ impl SetDisks { // Ok(()) // } - - // Shuffle the order - - // Shuffle the order - - // Return shuffled partsMetadata depending on distribution. - - // shuffle_disks TODO: use origin value } fn is_explicit_null_version(version_id: Option) -> bool { diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index a508b4027..1d9ceaf00 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -1181,8 +1181,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { checksum_combined.extend_from_slice(cs.raw.as_slice()); } - // TODO: check min part size - object_size += ext_part.size; object_actual_size += ext_part.actual_size; diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 553ece204..e1059d730 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -783,8 +783,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { .try_get_index() .map(crate::io_support::rio::compression_index_storage_bytes); - //TODO: userDefined - let mut etag = data.stream.try_resolve_etag().unwrap_or_default(); if let Some(ref tag) = opts.preserve_etag { etag = tag.clone(); @@ -1726,8 +1724,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { self.invalidate_get_object_metadata_cache(bucket, object).await; - // TODO: nslock - // Guard lock for metadata update let _lock_guard = if !opts.no_lock { Some(self.acquire_write_lock_diag("put_object_metadata", bucket, object).await?) @@ -2179,11 +2175,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { fi.metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags.to_owned()); - // TODO: userdeefined - self.update_object_meta(bucket, object, fi.clone(), disks.as_slice()).await?; - // TODO: versioned Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended)) } diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index f4d947255..5692f2364 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -1542,12 +1542,6 @@ impl ECStore { .instrument(tracing::Span::current()), ); - // let merge_res = merge_entry_channels(rx, inputs, sender.clone(), 1).await; - - // TODO: cancelList - - // let merge_res = merge_entry_channels(rx, inputs, sender.clone(), 1).await; - let results = join_all(futures).await; let mut all_at_eof = true; @@ -1572,10 +1566,6 @@ impl ECStore { return Ok(Vec::new()); } - // merge_res?; - - // TODO check cancel - for err in errs.iter() { if let Some(err) = err { if err == &Error::Unexpected { @@ -2010,9 +2000,6 @@ async fn gather_results( entry.name = entry.name.replace("\\", "/"); } - // TODO: rx.recv() - - // TODO: isLatestDeletemarker if !opts.include_directories && (entry.is_dir() || (!opts.versioned && entry.is_object() && entry.is_latest_delete_marker())) { @@ -2212,7 +2199,6 @@ async fn merge_entry_channels( if !dir_matches { // dir and object has the save name if other_entry.is_dir() { - // TODO: read next entry to top if !select_from(&rx, &mut in_channels, other_idx, &mut top, &mut n_done).await? { return Ok(()); } diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 2a46a2dac..414ac7cde 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -407,7 +407,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for ECStore { self.handle_verify_object_integrity(bucket, object, opts).await } - // TODO: review #[instrument(skip(self))] async fn copy_object( &self, @@ -436,7 +435,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for ECStore { async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { self.handle_delete_object(bucket, object, opts).await } - // TODO: review + #[instrument(skip(self))] async fn delete_objects( &self, diff --git a/crates/filemeta/src/fileinfo.rs b/crates/filemeta/src/fileinfo.rs index 679048599..cf562be8d 100644 --- a/crates/filemeta/src/fileinfo.rs +++ b/crates/filemeta/src/fileinfo.rs @@ -503,9 +503,7 @@ impl FileInfo { /// Check if replication related fields are equal pub fn replication_info_equals(&self, other: &FileInfo) -> bool { - self.mark_deleted == other.mark_deleted - // TODO: Add replication_state comparison when implemented - // && self.replication_state == other.replication_state + self.mark_deleted == other.mark_deleted && self.replication_state_internal == other.replication_state_internal } pub fn version_purge_status(&self) -> VersionPurgeStatusType { @@ -889,4 +887,31 @@ mod tests { prop_assert!(result.is_ok(), "FileInfo::unmarshal panicked for arbitrary input"); } } + + #[test] + fn replication_info_equals_compares_mark_deleted_and_replication_state() { + let base = FileInfo::default(); + + // Identical (both default) infos are equal. + assert!(base.replication_info_equals(&FileInfo::default())); + + // Differing mark_deleted breaks equality. + let mut marked = FileInfo::default(); + marked.mark_deleted = true; + assert!(!base.replication_info_equals(&marked)); + + // Differing replication_state_internal breaks equality (regression guard: + // this field used to be ignored by replication_info_equals). + let mut with_state = FileInfo::default(); + with_state.replication_state_internal = Some(ReplicationState { + replicate_decision_str: "arn:aws:s3:::dest".to_string(), + ..Default::default() + }); + assert!(!base.replication_info_equals(&with_state)); + + // Equal replication states remain equal. + let mut with_state_clone = FileInfo::default(); + with_state_clone.replication_state_internal = with_state.replication_state_internal.clone(); + assert!(with_state.replication_info_equals(&with_state_clone)); + } } diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index 27bded18c..9102a5553 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -111,7 +111,7 @@ pub use version::*; #[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)] pub struct FileMeta { pub versions: Vec, - pub data: InlineData, // TODO: xlMetaInlineData + pub data: InlineData, pub meta_ver: u8, } @@ -707,7 +707,6 @@ impl FileMeta { for ver in self.versions.iter() { let header = &ver.header; - // TODO: freeVersion if header.free_version() { non_free_versions -= 1; if include_free_versions diff --git a/crates/iam/src/error.rs b/crates/iam/src/error.rs index bf0cc9469..d434eb28a 100644 --- a/crates/iam/src/error.rs +++ b/crates/iam/src/error.rs @@ -286,21 +286,6 @@ pub fn is_err_no_such_service_account(err: &Error) -> bool { matches!(err, Error::NoSuchServiceAccount(_)) } -// pub fn clone_err(e: &Error) -> Error { -// if let Some(e) = e.downcast_ref::() { -// clone_disk_err(e) -// } else if let Some(e) = e.downcast_ref::() { -// if let Some(code) = e.raw_os_error() { -// Error::new(std::io::Error::from_raw_os_error(code)) -// } else { -// Error::new(std::io::Error::new(e.kind(), e.to_string())) -// } -// } else { -// //TODO: Optimize other types -// Error::msg(e.to_string()) -// } -// } - #[cfg(test)] mod tests { use super::*; diff --git a/crates/iam/src/manager.rs b/crates/iam/src/manager.rs index e84d6b1d2..7ba030762 100644 --- a/crates/iam/src/manager.rs +++ b/crates/iam/src/manager.rs @@ -384,7 +384,6 @@ where Ok(()) } - // TODO: Check if exists, whether retry is possible #[tracing::instrument(level = "debug", skip(self))] async fn save_iam_formatter(self: Arc) -> Result<()> { let path = get_iam_format_file_path(); diff --git a/crates/iam/src/store/object.rs b/crates/iam/src/store/object.rs index fad01e10f..d41834554 100644 --- a/crates/iam/src/store/object.rs +++ b/crates/iam/src/store/object.rs @@ -355,10 +355,6 @@ impl ObjectStore { async fn list_iam_config_items(&self, prefix: &str, ctx: CancellationToken, sender: Sender) { // debug!("list iam config items, prefix: {}", &prefix); - // TODO: Implement walk, use walk - - // let prefix = format!("{}{}", prefix, item); - let store = self.object_api.clone(); let (tx, mut rx) = mpsc::channel::(100); @@ -595,42 +591,6 @@ impl ObjectStore { Err(e) => Err(e.into()), } } - - // async fn load_policy(&self, name: &str) -> Result { - // let mut policy = self - // .load_iam_config::(&format!("config/iam/policies/{name}/policy.json")) - // .await?; - - // // FIXME: - // // if policy.version == 0 { - // // policy.create_date = object.mod_time; - // // policy.update_date = object.mod_time; - // // } - - // Ok(policy) - // } - - // async fn load_user_identity(&self, user_type: UserType, name: &str) -> Result> { - // let mut user = self - // .load_iam_config::(&format!( - // "config/iam/{base}{name}/identity.json", - // base = user_type.prefix(), - // name = name - // )) - // .await?; - - // if user.credentials.is_expired() { - // return Ok(None); - // } - - // if user.credentials.access_key.is_empty() { - // user.credentials.access_key = name.to_owned(); - // } - - // // todo, validate session token - - // Ok(Some(user)) - // } } #[async_trait::async_trait] diff --git a/crates/s3select-api/src/query/execution.rs b/crates/s3select-api/src/query/execution.rs index 7e4aa62d6..aa350c5f2 100644 --- a/crates/s3select-api/src/query/execution.rs +++ b/crates/s3select-api/src/query/execution.rs @@ -228,17 +228,14 @@ impl QueryStateMachine { } pub fn finish(&self) { - // TODO self.translate_to(QueryState::DONE(DONE::FINISHED)); } pub fn cancel(&self) { - // TODO self.translate_to(QueryState::DONE(DONE::CANCELLED)); } pub fn fail(&self) { - // TODO self.translate_to(QueryState::DONE(DONE::FAILED)); } diff --git a/crates/s3select-query/src/dispatcher/manager.rs b/crates/s3select-query/src/dispatcher/manager.rs index b18daad49..46a145204 100644 --- a/crates/s3select-query/src/dispatcher/manager.rs +++ b/crates/s3select-query/src/dispatcher/manager.rs @@ -195,8 +195,6 @@ impl SimpleQueryDispatcher { { file_format = file_format.with_delimiter(delimiter.as_bytes()[0]); } - // TODO waiting for processing @junxiang Mu - // if csv.file_header_info.is_some() {} match csv.file_header_info.as_ref() { Some(info) => { if *info == *NONE { diff --git a/crates/s3select-query/src/metadata/mod.rs b/crates/s3select-query/src/metadata/mod.rs index 2c8eb3a19..a682988d8 100644 --- a/crates/s3select-query/src/metadata/mod.rs +++ b/crates/s3select-query/src/metadata/mod.rs @@ -115,7 +115,6 @@ impl ContextProvider for MetadataProvider { } fn options(&self) -> &ConfigOptions { - // TODO refactor &self.config_options } diff --git a/crates/s3select-query/src/sql/logical/optimizer.rs b/crates/s3select-query/src/sql/logical/optimizer.rs index b1c669321..b0a27e4b5 100644 --- a/crates/s3select-query/src/sql/logical/optimizer.rs +++ b/crates/s3select-query/src/sql/logical/optimizer.rs @@ -43,8 +43,6 @@ pub trait LogicalOptimizer: Send + Sync { } pub struct DefaultLogicalOptimizer { - // fit datafusion - // TODO refactor analyzer: AnalyzerRef, rules: Vec>, } diff --git a/crates/signer/src/request_signature_v4.rs b/crates/signer/src/request_signature_v4.rs index 9735aa458..b93fbd927 100644 --- a/crates/signer/src/request_signature_v4.rs +++ b/crates/signer/src/request_signature_v4.rs @@ -62,8 +62,7 @@ struct SignFailure { type SignOutcome = std::result::Result, Box>; -#[allow(non_upper_case_globals)] // FIXME -static v4_ignored_headers: LazyLock> = LazyLock::new(|| { +static V4_IGNORED_HEADERS: LazyLock> = LazyLock::new(|| { let mut s = HashSet::new(); s.insert("accept-encoding"); s.insert("authorization"); @@ -307,7 +306,7 @@ fn pre_sign_v4_inner( } let credential = get_credential(access_key_id, location, t, SERVICE_TYPE_S3); - let signed_headers = get_signed_headers(&req, &v4_ignored_headers); + let signed_headers = get_signed_headers(&req, &V4_IGNORED_HEADERS); let mut query = >::new(); if let Some(q) = req.uri().query() { @@ -353,7 +352,7 @@ fn pre_sign_v4_inner( Ok(value) => value, Err(err) => return fail(req, err), }; - let canonical_request = match try_get_canonical_request(&req, &v4_ignored_headers, &hashed_payload) { + let canonical_request = match try_get_canonical_request(&req, &V4_IGNORED_HEADERS, &hashed_payload) { Ok(value) => value, Err(err) => return fail(req, err), }; @@ -557,7 +556,7 @@ fn sign_v4_inner( Ok(value) => value, Err(err) => return fail(req, err), }; - let canonical_request = match try_get_canonical_request(&req, &v4_ignored_headers, &hashed_payload) { + let canonical_request = match try_get_canonical_request(&req, &V4_IGNORED_HEADERS, &hashed_payload) { Ok(value) => value, Err(err) => return fail(req, err), }; @@ -567,7 +566,7 @@ fn sign_v4_inner( }; let signing_key = get_signing_key(secret_access_key, location, t, service_type); let credential = get_credential(access_key_id, location, t2, service_type); - let signed_headers = get_signed_headers(&req, &v4_ignored_headers); + let signed_headers = get_signed_headers(&req, &V4_IGNORED_HEADERS); let signature = get_signature(signing_key, &string_to_sign); //debug!("\n\ncanonical_request: \n{}\nstring_to_sign: \n{}\nsignature: \n{}\n\n", &canonical_request, &string_to_sign, &signature); @@ -742,7 +741,7 @@ mod tests { let hashed_payload = try_get_hashed_payload(&req).expect("example request should have valid payload header"); let canonical_request = - try_get_canonical_request(&req, &v4_ignored_headers, &hashed_payload).expect("example request should canonicalize"); + try_get_canonical_request(&req, &V4_IGNORED_HEADERS, &hashed_payload).expect("example request should canonicalize"); assert_eq!( canonical_request, concat!( @@ -819,7 +818,7 @@ mod tests { let hashed_payload = try_get_hashed_payload(&req).expect("example request should have valid payload header"); let canonical_request = - try_get_canonical_request(&req, &v4_ignored_headers, &hashed_payload).expect("example request should canonicalize"); + try_get_canonical_request(&req, &V4_IGNORED_HEADERS, &hashed_payload).expect("example request should canonicalize"); println!("canonical_request: \n{canonical_request}\n"); assert_eq!( canonical_request, @@ -887,7 +886,7 @@ mod tests { println!("{:?}", req.uri().query()); let hashed_payload = try_get_hashed_payload(&req).expect("example request should have valid payload header"); let canonical_request = - try_get_canonical_request(&req, &v4_ignored_headers, &hashed_payload).expect("example request should canonicalize"); + try_get_canonical_request(&req, &V4_IGNORED_HEADERS, &hashed_payload).expect("example request should canonicalize"); println!("canonical_request: \n{canonical_request}\n"); assert_eq!( canonical_request, @@ -955,7 +954,7 @@ mod tests { println!("{:?}", req.uri().query()); let hashed_payload = try_get_hashed_payload(&req).expect("example request should have valid payload header"); let canonical_request = - try_get_canonical_request(&req, &v4_ignored_headers, &hashed_payload).expect("example request should canonicalize"); + try_get_canonical_request(&req, &V4_IGNORED_HEADERS, &hashed_payload).expect("example request should canonicalize"); println!("canonical_request: \n{canonical_request}\n"); assert_eq!( canonical_request, @@ -1024,12 +1023,12 @@ mod tests { canonical_request.push_str(req.uri().query().unwrap()); canonical_request.push('\n'); canonical_request.push_str( - try_get_canonical_headers(&req, &v4_ignored_headers) + try_get_canonical_headers(&req, &V4_IGNORED_HEADERS) .expect("presigned request should canonicalize headers") .as_str(), ); canonical_request.push('\n'); - canonical_request.push_str(&get_signed_headers(&req, &v4_ignored_headers)); + canonical_request.push_str(&get_signed_headers(&req, &V4_IGNORED_HEADERS)); canonical_request.push('\n'); canonical_request.push_str( try_get_hashed_payload(&req) @@ -1080,12 +1079,12 @@ mod tests { canonical_request.push_str(req.uri().query().unwrap()); canonical_request.push('\n'); canonical_request.push_str( - try_get_canonical_headers(&req, &v4_ignored_headers) + try_get_canonical_headers(&req, &V4_IGNORED_HEADERS) .expect("presigned request should canonicalize headers") .as_str(), ); canonical_request.push('\n'); - canonical_request.push_str(&get_signed_headers(&req, &v4_ignored_headers)); + canonical_request.push_str(&get_signed_headers(&req, &V4_IGNORED_HEADERS)); canonical_request.push('\n'); canonical_request.push_str( try_get_hashed_payload(&req) diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index 35111333c..cfefbe0dc 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -1153,12 +1153,21 @@ mod tests { assert!(err.to_string().contains("invalid bucket name")); } - #[ignore] // FIXME: failed in github actions - keeping original test #[test] fn test_decode() { + // Mirror the production decode path (handler decodes the request body via + // serde_json::from_slice), not urlencoded. let b = b"{\"recursive\":false,\"dryRun\":false,\"remove\":false,\"recreate\":false,\"scanMode\":1,\"updateParity\":false,\"nolock\":false}"; - let s: HealOpts = serde_urlencoded::from_bytes(b).unwrap(); - debug!("Parsed HealOpts: {:?}", s); + let s: HealOpts = serde_json::from_slice(b).expect("HealOpts JSON body must decode"); + assert!(!s.recursive); + assert!(!s.dry_run); + assert!(!s.remove); + assert!(!s.recreate); + assert_eq!(s.scan_mode, HealScanMode::Normal); + assert!(!s.update_parity); + assert!(!s.no_lock); + assert_eq!(s.pool, None); + assert_eq!(s.set, None); } #[tokio::test] diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 931b65a99..c8588686a 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -5398,8 +5398,6 @@ impl DefaultObjectUsecase { }; let last_modified = info.mod_time.map(Timestamp::from); - // TODO: range download - let content_length = info.get_actual_size().map_err(|e| { error!(error = %e, "Failed to resolve actual object size"); ApiError::from(e) diff --git a/rustfs/src/auth.rs b/rustfs/src/auth.rs index 1aff4172b..458575022 100644 --- a/rustfs/src/auth.rs +++ b/rustfs/src/auth.rs @@ -439,8 +439,6 @@ pub fn check_claims_from_token(token: &str, cred: &Credentials) -> S3Result