diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 6ca6978d2..425f35d83 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -588,6 +588,18 @@ impl SetDisks { // shuffle_disks TODO: use origin value } +fn is_explicit_null_version(version_id: Option) -> bool { + version_id == Some(Uuid::nil()) +} + +fn delete_file_info_version_id(version_id: Option) -> Option { + if is_explicit_null_version(version_id) { + None + } else { + version_id + } +} + #[async_trait::async_trait] impl ObjectIO for SetDisks { #[tracing::instrument(level = "debug", skip(self))] @@ -1606,9 +1618,10 @@ impl ObjectOperations for SetDisks { let mut vers_map: HashMap<&String, FileInfoVersions> = HashMap::new(); for (i, dobj) in objects.iter().enumerate() { + let explicit_null_version = is_explicit_null_version(dobj.version_id); let mut vr = FileInfo { name: dobj.object_name.clone(), - version_id: dobj.version_id, + version_id: delete_file_info_version_id(dobj.version_id), idx: i, replication_state_internal: Some(dobj.replication_state()), ..Default::default() @@ -1657,7 +1670,11 @@ impl ObjectOperations for SetDisks { } else { del_objects[i] = DeletedObject { object_name: vr.name.clone(), - version_id: vr.version_id, + version_id: if explicit_null_version { + Some(Uuid::nil()) + } else { + vr.version_id + }, replication_state: vr.replication_state_internal.clone(), ..Default::default() } @@ -6022,6 +6039,15 @@ mod tests { assert!(parts_after_marker(&part_numbers, 4).is_none()); } + #[test] + fn delete_file_info_version_id_maps_explicit_null_version_to_stored_null() { + assert_eq!(delete_file_info_version_id(Some(Uuid::nil())), None); + + let version_id = Uuid::new_v4(); + assert_eq!(delete_file_info_version_id(Some(version_id)), Some(version_id)); + assert_eq!(delete_file_info_version_id(None), None); + } + #[test] fn test_is_cold_storage_class() { // Test cold storage classes diff --git a/crates/ecstore/src/store_api/types.rs b/crates/ecstore/src/store_api/types.rs index 42a2c7f6f..4f43dde66 100644 --- a/crates/ecstore/src/store_api/types.rs +++ b/crates/ecstore/src/store_api/types.rs @@ -615,11 +615,12 @@ impl ObjectInfo { bucket: &str, prefix: &str, delimiter: Option, - after_version_id: Option, + after_version_marker: Option, ) -> Vec { let vcfg = get_versioning_config(bucket).await.ok(); let mut objects = Vec::with_capacity(entries.entries().len()); let mut prev_prefix = ""; + let mut after_version_marker = after_version_marker; for entry in entries.entries() { if entry.is_object() { if let Some(delimiter) = &delimiter { @@ -656,12 +657,8 @@ impl ObjectInfo { } }; - let versions = if let Some(vid) = after_version_id { - if let Some(idx) = file_infos.find_version_index(vid) { - &file_infos.versions[idx + 1..] - } else { - &file_infos.versions - } + let versions = if let Some(marker) = after_version_marker.take() { + versions_after_marker(&file_infos, marker) } else { &file_infos.versions }; @@ -833,6 +830,23 @@ impl ObjectInfo { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum VersionMarker { + Null, + Version(Uuid), +} + +fn versions_after_marker(file_infos: &rustfs_filemeta::FileInfoVersions, marker: VersionMarker) -> &[FileInfo] { + let marker_idx = match marker { + VersionMarker::Null => file_infos.versions.iter().position(|version| version.version_id.is_none()), + VersionMarker::Version(vid) => file_infos.find_version_index(vid), + }; + + marker_idx + .map(|idx| &file_infos.versions[idx + 1..]) + .unwrap_or(&file_infos.versions) +} + #[derive(Debug, Default)] pub struct ListObjectsInfo { // Indicates whether the returned list objects response is truncated. A @@ -1076,6 +1090,100 @@ mod tests { use super::*; use rustfs_filemeta::ReplicationState; + #[test] + fn versions_after_marker_handles_null_version_marker() { + let first_version = Uuid::parse_str("11111111-2222-3333-4444-555555555555").unwrap(); + let last_version = Uuid::parse_str("aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee").unwrap(); + let file_infos = rustfs_filemeta::FileInfoVersions { + versions: vec![ + FileInfo { + version_id: Some(first_version), + ..Default::default() + }, + FileInfo { + version_id: None, + ..Default::default() + }, + FileInfo { + version_id: Some(last_version), + ..Default::default() + }, + ], + ..Default::default() + }; + + let versions = versions_after_marker(&file_infos, VersionMarker::Null); + + assert_eq!(versions.len(), 1); + assert_eq!(versions[0].version_id, Some(last_version)); + } + + #[test] + fn versions_after_marker_handles_uuid_version_marker() { + let first_version = Uuid::parse_str("11111111-2222-3333-4444-555555555555").unwrap(); + let last_version = Uuid::parse_str("aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee").unwrap(); + let file_infos = rustfs_filemeta::FileInfoVersions { + versions: vec![ + FileInfo { + version_id: Some(first_version), + ..Default::default() + }, + FileInfo { + version_id: None, + ..Default::default() + }, + FileInfo { + version_id: Some(last_version), + ..Default::default() + }, + ], + ..Default::default() + }; + + let versions = versions_after_marker(&file_infos, VersionMarker::Version(first_version)); + + assert_eq!(versions.len(), 2); + assert_eq!(versions[0].version_id, None); + assert_eq!(versions[1].version_id, Some(last_version)); + } + + #[tokio::test] + async fn versions_listing_applies_version_marker_only_to_first_entry() { + let metadata = rustfs_filemeta::test_data::create_real_xlmeta().expect("test metadata should be valid"); + let entries = rustfs_filemeta::MetaCacheEntriesSorted { + o: rustfs_filemeta::MetaCacheEntries(vec![ + Some(rustfs_filemeta::MetaCacheEntry { + name: "obj-a".to_owned(), + metadata: metadata.clone(), + ..Default::default() + }), + Some(rustfs_filemeta::MetaCacheEntry { + name: "obj-b".to_owned(), + metadata, + ..Default::default() + }), + ]), + ..Default::default() + }; + let marker_version = Uuid::parse_str("11111111-2222-3333-4444-555555555555").unwrap(); + + let objects = ObjectInfo::from_meta_cache_entries_sorted_versions( + &entries, + "bucket", + "", + None, + Some(VersionMarker::Version(marker_version)), + ) + .await; + + let obj_a_count = objects.iter().filter(|object| object.name == "obj-a").count(); + let obj_b_count = objects.iter().filter(|object| object.name == "obj-b").count(); + + assert_eq!(obj_a_count, 2); + assert_eq!(obj_b_count, 3); + assert_eq!(objects.len(), 5); + } + #[test] fn get_actual_size_prefers_actual_size_field() { let info = ObjectInfo { diff --git a/crates/ecstore/src/store_list_objects.rs b/crates/ecstore/src/store_list_objects.rs index dd3d3b38a..20332366c 100644 --- a/crates/ecstore/src/store_list_objects.rs +++ b/crates/ecstore/src/store_list_objects.rs @@ -23,8 +23,8 @@ use crate::error::{ }; use crate::set_disk::SetDisks; use crate::store_api::{ - ListObjectVersionsInfo, ListObjectsInfo, ObjectInfo, ObjectInfoOrErr, ObjectOperations, ObjectOptions, WalkOptions, - WalkVersionsSortOrder, + ListObjectVersionsInfo, ListObjectsInfo, ObjectInfo, ObjectInfoOrErr, ObjectOperations, ObjectOptions, VersionMarker, + WalkOptions, WalkVersionsSortOrder, }; use crate::store_utils::is_reserved_or_invalid_bucket; use crate::{store::ECStore, store_api::ListObjectsV2Info}; @@ -118,9 +118,12 @@ pub struct ListPathOptions { pub filter_prefix: Option, // Marker to resume listing. - // The response will be the first entry >= this object name. pub marker: Option, + // Include marker itself in the returned entries. Version listings need this + // when a version marker selects a later version from the marker object. + pub include_marker: bool, + // Limit the number of results. pub limit: i32, @@ -159,6 +162,67 @@ pub struct ListPathOptions { } const MARKER_TAG_VERSION: &str = "v1"; +const ENV_API_LIST_QUORUM: &str = "RUSTFS_API_LIST_QUORUM"; +const DEFAULT_API_LIST_QUORUM: &str = "strict"; + +fn normalize_list_quorum(value: &str) -> &'static str { + let value = value.trim(); + if value.eq_ignore_ascii_case("disk") { + "disk" + } else if value.eq_ignore_ascii_case("reduced") { + "reduced" + } else if value.eq_ignore_ascii_case("optimal") { + "optimal" + } else if value.eq_ignore_ascii_case("auto") { + "auto" + } else { + DEFAULT_API_LIST_QUORUM + } +} + +fn list_quorum_from_env() -> String { + let value = rustfs_utils::get_env_str(ENV_API_LIST_QUORUM, DEFAULT_API_LIST_QUORUM); + normalize_list_quorum(&value).to_owned() +} + +fn list_metadata_resolution_params(bucket: String, listing_quorum: usize, versioned: bool) -> MetadataResolutionParams { + let mut resolver = MetadataResolutionParams { + dir_quorum: listing_quorum, + obj_quorum: listing_quorum, + bucket, + ..Default::default() + }; + + if !versioned { + resolver.requested_versions = 1; + } + + resolver +} + +fn parse_version_marker(marker: String) -> Result { + if marker == "null" { + Ok(VersionMarker::Null) + } else { + Ok(VersionMarker::Version(Uuid::parse_str(&marker)?)) + } +} + +fn version_marker_for_entries( + entries: Option<&MetaCacheEntriesSorted>, + key_marker: Option<&str>, + version_marker: Option, +) -> Option { + let marker = version_marker?; + let Some(key_marker) = key_marker else { + return Some(marker); + }; + + entries + .and_then(|entries| entries.entries().first().map(|entry| entry.name.as_str() == key_marker)) + .unwrap_or_default() + .then_some(marker) +} impl ListPathOptions { pub fn set_filter(&mut self) { @@ -182,50 +246,64 @@ impl ListPathOptions { } pub fn parse_marker(&mut self) { - if let Some(marker) = &self.marker { - let s = marker.clone(); - if !s.contains(format!("[rustfs_cache:{MARKER_TAG_VERSION}").as_str()) { - return; - } + let Some(marker) = self.marker.clone() else { + return; + }; + let Some(start_idx) = marker.rfind("[rustfs_cache:") else { + return; + }; + let Some(end_offset) = marker[start_idx..].rfind(']') else { + return; + }; - if let (Some(start_idx), Some(end_idx)) = (s.find("["), s.find("]")) { - self.marker = Some(s[0..start_idx].to_owned()); - let tags: Vec<_> = s[start_idx..end_idx].trim_matches(['[', ']']).split(",").collect(); + let end_idx = start_idx + end_offset; + let tag_body = marker[start_idx + 1..end_idx].to_owned(); + let mut supported_marker = false; - for &tag in tags.iter() { - let kv: Vec<_> = tag.split(":").collect(); - if kv.len() != 2 { - continue; - } - - match kv[0] { - "rustfs_cache" if kv[1] != MARKER_TAG_VERSION => { - continue; - } - "id" => self.id = Some(kv[1].to_owned()), - "return" => { - self.id = Some(Uuid::new_v4().to_string()); - self.create = true; - } - "p" => match kv[1].parse::() { - Ok(res) => self.pool_idx = Some(res), - Err(_) => { - self.id = Some(Uuid::new_v4().to_string()); - self.create = true; - continue; - } - }, - "s" => match kv[1].parse::() { - Ok(res) => self.set_idx = Some(res), - Err(_) => { - self.id = Some(Uuid::new_v4().to_string()); - self.create = true; - continue; - } - }, - _ => (), - } + for tag in tag_body.split(',') { + let Some((key, value)) = tag.split_once(':') else { + continue; + }; + if key == "rustfs_cache" { + if value != MARKER_TAG_VERSION { + return; } + supported_marker = true; + } + } + + if !supported_marker { + return; + } + + self.marker = Some(marker[..start_idx].to_owned()); + for tag in tag_body.split(',') { + let Some((key, value)) = tag.split_once(':') else { + continue; + }; + + match key { + "rustfs_cache" => {} + "id" => self.id = Some(value.to_owned()), + "return" => { + self.id = Some(Uuid::new_v4().to_string()); + self.create = true; + } + "p" => match value.parse::() { + Ok(res) => self.pool_idx = Some(res), + Err(_) => { + self.id = Some(Uuid::new_v4().to_string()); + self.create = true; + } + }, + "s" => match value.parse::() { + Ok(res) => self.set_idx = Some(res), + Err(_) => { + self.id = Some(Uuid::new_v4().to_string()); + self.create = true; + } + }, + _ => (), } } } @@ -301,7 +379,7 @@ impl ECStore { limit: effective_max_keys, marker, incl_deleted, - ask_disks: "strict".to_owned(), //TODO: from config + ask_disks: list_quorum_from_env(), ..Default::default() }; @@ -444,13 +522,9 @@ impl ECStore { return Err(StorageError::NotImplemented); } + let has_version_marker = version_marker.is_some(); let version_marker = if let Some(marker) = version_marker { - // "null" is used for non-versioned objects in AWS S3 API - if marker == "null" { - None - } else { - Some(Uuid::parse_str(&marker)?) - } + Some(parse_version_marker(marker)?) } else { None }; @@ -464,8 +538,9 @@ impl ECStore { limit: effective_max_keys, marker, incl_deleted: true, - ask_disks: "strict".to_owned(), + ask_disks: list_quorum_from_env(), versioned: true, + include_marker: has_version_marker, ..Default::default() }; @@ -486,10 +561,14 @@ impl ECStore { return Err(to_object_err(err.into(), vec![bucket, prefix])); } - if let Some(result) = list_result.entries.as_mut() { - result.forward_past(opts.marker); + if let Some(result) = list_result.entries.as_mut() + && !has_version_marker + { + result.forward_past(opts.marker.clone()); } + let version_marker = version_marker_for_entries(list_result.entries.as_ref(), opts.marker.as_deref(), version_marker); + let mut get_objects = ObjectInfo::from_meta_cache_entries_sorted_versions( &list_result.entries.unwrap_or_default(), bucket, @@ -1080,7 +1159,7 @@ async fn gather_results( } if let Some(marker) = &opts.marker - && &entry.name <= marker + && ((!opts.include_marker && &entry.name <= marker) || (opts.include_marker && &entry.name < marker)) { continue; } @@ -1349,16 +1428,7 @@ impl SetDisks { fallback_disks = disks.split_off(ask_disks as usize); } - let mut resolver = MetadataResolutionParams { - dir_quorum: listing_quorum, - obj_quorum: listing_quorum, - bucket: opts.bucket.clone(), - ..Default::default() - }; - - if opts.versioned { - resolver.requested_versions = 1; - } + let resolver = list_metadata_resolution_params(opts.bucket.clone(), listing_quorum, opts.versioned); let limit = { if opts.limit > 0 && opts.stop_disk_at_limit { @@ -1483,9 +1553,13 @@ fn calc_common_counter(infos: &[DiskInfo], read_quorum: usize) -> u64 { #[cfg(test)] mod test { - use super::{ListPathOptions, MAX_OBJECT_LIST, gather_results, max_keys_plus_one, walk_result_from_set_errors}; + use super::{ + ENV_API_LIST_QUORUM, ListPathOptions, MAX_OBJECT_LIST, VersionMarker, gather_results, list_metadata_resolution_params, + list_quorum_from_env, max_keys_plus_one, normalize_list_quorum, parse_version_marker, version_marker_for_entries, + walk_result_from_set_errors, + }; use crate::error::StorageError; - use rustfs_filemeta::MetaCacheEntry; + use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntriesSorted, MetaCacheEntry}; use std::time::Duration; use tokio::sync::mpsc; use tokio::time::timeout; @@ -1499,6 +1573,13 @@ mod test { } } + fn sorted_entries(names: &[&str]) -> MetaCacheEntriesSorted { + MetaCacheEntriesSorted { + o: MetaCacheEntries(names.iter().map(|name| Some(test_meta_entry(name))).collect()), + ..Default::default() + } + } + #[tokio::test] async fn gather_results_returns_after_limit_without_waiting_for_input_close() { let (entry_tx, entry_rx) = mpsc::channel(4); @@ -1531,6 +1612,109 @@ mod test { .expect("gather_results should succeed"); } + #[tokio::test] + async fn gather_results_keeps_marker_entry_for_version_marker_listing() { + let (entry_tx, entry_rx) = mpsc::channel(4); + let (result_tx, mut result_rx) = mpsc::channel(1); + + entry_tx.send(test_meta_entry("obj-a")).await.unwrap(); + entry_tx.send(test_meta_entry("obj-b")).await.unwrap(); + + let handle = tokio::spawn(gather_results( + CancellationToken::new(), + ListPathOptions { + bucket: "bucket".to_owned(), + marker: Some("obj-a".to_owned()), + include_marker: true, + limit: 2, + incl_deleted: true, + versioned: true, + ..Default::default() + }, + entry_rx, + result_tx, + )); + + let result = timeout(Duration::from_secs(1), result_rx.recv()) + .await + .expect("limited result should be sent promptly") + .expect("limited result should be present"); + let entries = result.entries.unwrap(); + let names = entries + .entries() + .into_iter() + .map(|entry| entry.name.as_str()) + .collect::>(); + + assert_eq!(names, ["obj-a", "obj-b"]); + + timeout(Duration::from_secs(1), handle) + .await + .expect("gather_results should finish after sending a limited result") + .expect("gather_results task should not panic") + .expect("gather_results should succeed"); + } + + #[test] + fn version_marker_is_applied_only_when_key_marker_entry_is_present() { + let version_marker = Some(VersionMarker::Null); + + let listed_after_deleted_marker = sorted_entries(&["obj-b", "obj-c"]); + assert_eq!( + version_marker_for_entries(Some(&listed_after_deleted_marker), Some("obj-a"), version_marker), + None + ); + + let listed_with_marker = sorted_entries(&["obj-a", "obj-b"]); + assert_eq!( + version_marker_for_entries(Some(&listed_with_marker), Some("obj-a"), version_marker), + version_marker + ); + } + + #[tokio::test] + async fn gather_results_skips_marker_entry_by_default() { + let (entry_tx, entry_rx) = mpsc::channel(4); + let (result_tx, mut result_rx) = mpsc::channel(1); + + entry_tx.send(test_meta_entry("obj-a")).await.unwrap(); + entry_tx.send(test_meta_entry("obj-b")).await.unwrap(); + drop(entry_tx); + + let handle = tokio::spawn(gather_results( + CancellationToken::new(), + ListPathOptions { + bucket: "bucket".to_owned(), + marker: Some("obj-a".to_owned()), + limit: 2, + incl_deleted: true, + ..Default::default() + }, + entry_rx, + result_tx, + )); + + let result = timeout(Duration::from_secs(1), result_rx.recv()) + .await + .expect("eof result should be sent promptly") + .expect("eof result should be present"); + let entries = result.entries.unwrap(); + let names = entries + .entries() + .into_iter() + .map(|entry| entry.name.as_str()) + .collect::>(); + + assert_eq!(names, ["obj-b"]); + assert!(result.err.is_some()); + + timeout(Duration::from_secs(1), handle) + .await + .expect("gather_results should finish after input closes") + .expect("gather_results task should not panic") + .expect("gather_results should succeed"); + } + #[test] fn test_max_keys_plus_one_caps_before_lookahead() { assert_eq!(max_keys_plus_one(999, true), 1000); @@ -1540,28 +1724,74 @@ mod test { assert_eq!(max_keys_plus_one(-1, true), MAX_OBJECT_LIST + 1); } - /// Test that "null" version marker is handled correctly - /// AWS S3 API uses "null" string to represent non-versioned objects + #[test] + fn normalize_list_quorum_accepts_supported_values() { + assert_eq!(normalize_list_quorum("strict"), "strict"); + assert_eq!(normalize_list_quorum("disk"), "disk"); + assert_eq!(normalize_list_quorum("reduced"), "reduced"); + assert_eq!(normalize_list_quorum("optimal"), "optimal"); + assert_eq!(normalize_list_quorum("auto"), "auto"); + assert_eq!(normalize_list_quorum(" OPTIMAL "), "optimal"); + } + + #[test] + fn normalize_list_quorum_falls_back_to_strict() { + assert_eq!(normalize_list_quorum(""), "strict"); + assert_eq!(normalize_list_quorum("unknown"), "strict"); + } + + #[test] + #[serial_test::serial] + fn list_quorum_from_env_defaults_to_strict() { + temp_env::with_var_unset(ENV_API_LIST_QUORUM, || { + assert_eq!(list_quorum_from_env(), "strict"); + }); + } + + #[test] + #[serial_test::serial] + fn list_quorum_from_env_honors_supported_value() { + temp_env::with_var(ENV_API_LIST_QUORUM, Some("auto"), || { + assert_eq!(list_quorum_from_env(), "auto"); + }); + } + + #[test] + #[serial_test::serial] + fn list_quorum_from_env_rejects_unknown_value() { + temp_env::with_var(ENV_API_LIST_QUORUM, Some("unsafe"), || { + assert_eq!(list_quorum_from_env(), "strict"); + }); + } + + #[test] + fn list_metadata_resolution_params_limits_plain_listing_to_latest_version() { + let resolver = list_metadata_resolution_params("bucket".to_string(), 2, false); + + assert_eq!(resolver.dir_quorum, 2); + assert_eq!(resolver.obj_quorum, 2); + assert_eq!(resolver.bucket, "bucket"); + assert_eq!(resolver.requested_versions, 1); + } + + #[test] + fn list_metadata_resolution_params_keeps_all_versions_for_version_listing() { + let resolver = list_metadata_resolution_params("bucket".to_string(), 3, true); + + assert_eq!(resolver.dir_quorum, 3); + assert_eq!(resolver.obj_quorum, 3); + assert_eq!(resolver.bucket, "bucket"); + assert_eq!(resolver.requested_versions, 0); + } + #[test] fn test_null_version_marker_handling() { - // "null" should be treated as None (non-versioned) - let version_marker = "null"; - let parsed: Option = if version_marker == "null" { - None - } else { - Uuid::parse_str(version_marker).ok() - }; - assert!(parsed.is_none(), "\"null\" should be parsed as None"); + let parsed = parse_version_marker("null".to_string()).expect("null marker should parse"); + assert_eq!(parsed, VersionMarker::Null); - // Valid UUID should be parsed correctly let valid_uuid = "550e8400-e29b-41d4-a716-446655440000"; - let parsed: Option = if valid_uuid == "null" { - None - } else { - Uuid::parse_str(valid_uuid).ok() - }; - assert!(parsed.is_some(), "Valid UUID should be parsed correctly"); - assert_eq!(parsed.unwrap().to_string(), "550e8400-e29b-41d4-a716-446655440000"); + let parsed = parse_version_marker(valid_uuid.to_string()).expect("uuid marker should parse"); + assert_eq!(parsed, VersionMarker::Version(Uuid::parse_str(valid_uuid).unwrap())); } /// Test that next_version_idmarker returns "null" for non-versioned objects @@ -1584,15 +1814,11 @@ mod test { // Scenario 1: Non-versioned object // Server returns "null" as NextVersionIdMarker // Client sends "null" as VersionIdMarker - // Server parses "null" as None + // Server parses "null" as an explicit null-version marker let server_response = "null"; let client_request = server_response; - let parsed: Option = if client_request == "null" { - None - } else { - Uuid::parse_str(client_request).ok() - }; - assert!(parsed.is_none()); + let parsed = parse_version_marker(client_request.to_string()).expect("null marker should parse"); + assert_eq!(parsed, VersionMarker::Null); // Scenario 2: Versioned object // Server returns UUID as NextVersionIdMarker @@ -1601,13 +1827,8 @@ mod test { let uuid_str = "550e8400-e29b-41d4-a716-446655440000"; let server_response = uuid_str; let client_request = server_response; - let parsed: Option = if client_request == "null" { - None - } else { - Uuid::parse_str(client_request).ok() - }; - assert!(parsed.is_some()); - assert_eq!(parsed.unwrap().to_string(), uuid_str); + let parsed = parse_version_marker(client_request.to_string()).expect("uuid marker should parse"); + assert_eq!(parsed, VersionMarker::Version(Uuid::parse_str(uuid_str).unwrap())); } #[test] @@ -1639,6 +1860,41 @@ mod test { assert!(!parsed.create); } + #[test] + fn list_path_marker_parser_uses_trailing_cache_tag() { + let mut parsed = ListPathOptions { + marker: Some(format!( + "photos/[archive]/image.jpg[rustfs_cache:{},id:list-cache-id,p:3,s:7]", + super::MARKER_TAG_VERSION + )), + ..Default::default() + }; + + parsed.parse_marker(); + + assert_eq!(parsed.marker.as_deref(), Some("photos/[archive]/image.jpg")); + assert_eq!(parsed.id.as_deref(), Some("list-cache-id")); + assert_eq!(parsed.pool_idx, Some(3)); + assert_eq!(parsed.set_idx, Some(7)); + } + + #[test] + fn list_path_marker_parser_ignores_unsupported_cache_tag_version() { + let marker = "photos/image.jpg[rustfs_cache:v0,id:list-cache-id,p:3,s:7]".to_string(); + let mut parsed = ListPathOptions { + marker: Some(marker.clone()), + ..Default::default() + }; + + parsed.parse_marker(); + + assert_eq!(parsed.marker.as_deref(), Some(marker.as_str())); + assert!(parsed.id.is_none()); + assert!(parsed.pool_idx.is_none()); + assert!(parsed.set_idx.is_none()); + assert!(!parsed.create); + } + #[test] fn walk_result_from_set_errors_returns_non_eof_error() { let err = walk_result_from_set_errors(&[Some(StorageError::Unexpected), Some(StorageError::FileAccessDenied)]) diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index ae750ae0f..a2edd9c9a 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -2022,6 +2022,7 @@ impl DefaultBucketUsecase { let ListObjectVersionsInput { bucket, delimiter, + encoding_type, key_marker, version_id_marker, max_keys, @@ -2029,22 +2030,23 @@ impl DefaultBucketUsecase { .. } = req.input; - let ListObjectVersionsParams { - prefix, - delimiter, - key_marker, - version_id_marker, - max_keys, - } = parse_list_object_versions_params(prefix, delimiter, key_marker, version_id_marker, max_keys)?; + let params = parse_list_object_versions_params(prefix, delimiter, key_marker, version_id_marker, max_keys)?; let store = get_validated_store(&bucket).await?; let object_infos = store - .list_object_versions(&bucket, &prefix, key_marker, version_id_marker, delimiter.clone(), max_keys) + .list_object_versions( + &bucket, + ¶ms.prefix, + params.key_marker.clone(), + params.version_id_marker.clone(), + params.delimiter.clone(), + params.max_keys, + ) .await .map_err(ApiError::from)?; - let output = build_list_object_versions_output(object_infos, bucket, prefix, delimiter, max_keys); + let output = build_list_object_versions_output(object_infos, bucket, ¶ms, encoding_type.as_ref()); Ok(S3Response::new(output)) } diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index bbab2badc..fa7eee0a9 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -3158,6 +3158,8 @@ impl DefaultObjectUsecase { key: Some(v.object_name.clone()), version_id: if is_dir_object(v.object_name.as_str()) && v.version_id == Some(Uuid::nil()) { None + } else if v.version_id == Some(Uuid::nil()) { + Some("null".to_string()) } else { v.version_id.map(|v| v.to_string()) }, @@ -5093,6 +5095,15 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[test] + fn normalize_delete_objects_version_id_preserves_explicit_null_marker() { + let (wire_version_id, internal_version_id) = + normalize_delete_objects_version_id(Some("null".to_string())).expect("null version marker should parse"); + + assert_eq!(wire_version_id.as_deref(), Some("null")); + assert_eq!(internal_version_id, Some(Uuid::nil())); + } + #[test] fn should_schedule_delete_replication_skips_replica_requests() { let opts = ObjectOptions { diff --git a/rustfs/src/storage/s3_api/bucket.rs b/rustfs/src/storage/s3_api/bucket.rs index c40ca405f..626b5bed0 100644 --- a/rustfs/src/storage/s3_api/bucket.rs +++ b/rustfs/src/storage/s3_api/bucket.rs @@ -13,6 +13,7 @@ // limitations under the License. use crate::storage::s3_api::common::rustfs_owner; +use percent_encoding::percent_decode_str; use rustfs_ecstore::client::object_api_utils::to_s3s_etag; use rustfs_ecstore::store_api::{BucketInfo, ListObjectVersionsInfo, ListObjectsV2Info}; use s3s::dto::{ @@ -29,6 +30,25 @@ fn normalize_max_keys(max_keys: i32) -> i32 { max_keys.min(S3_MAX_KEYS) } +fn should_encode_url(encoding_type: Option<&EncodingType>) -> bool { + encoding_type.is_some_and(|e| e.as_str() == EncodingType::URL) +} + +fn encode_s3_name(name: &str) -> String { + name.split('/') + .map(|part| encode(part).to_string()) + .collect::>() + .join("/") +} + +fn encode_list_output_value(value: String, encoding_type: Option<&EncodingType>) -> String { + if should_encode_url(encoding_type) { + encode_s3_name(&value) + } else { + value + } +} + #[derive(Debug, PartialEq, Eq)] pub(crate) struct ListObjectVersionsParams { pub prefix: String, @@ -140,12 +160,11 @@ pub(crate) fn parse_list_objects_v2_params( let start_after_for_query = start_after.filter(|v| !v.is_empty()); // Save original continuation_token for response (per S3 API spec, must echo back if provided) - // Note: empty string should still be echoed back in the response let response_continuation_token = continuation_token.clone(); - let continuation_token_for_query = continuation_token.filter(|v| !v.is_empty()); // Decode continuation_token from base64 for internal use - let decoded_continuation_token = continuation_token_for_query + let decoded_continuation_token = continuation_token + .filter(|token| !token.is_empty()) .map(|token| { base64_simd::STANDARD .decode_to_vec(token.as_bytes()) @@ -172,16 +191,16 @@ pub(crate) fn parse_list_objects_v2_params( pub(crate) fn build_list_object_versions_output( object_infos: ListObjectVersionsInfo, bucket: String, - prefix: String, - delimiter: Option, - max_keys: i32, + params: &ListObjectVersionsParams, + encoding_type: Option<&EncodingType>, ) -> ListObjectVersionsOutput { + let encode_output_value = |value: &str| encode_list_output_value(value.to_owned(), encoding_type); let versions: Vec = object_infos .objects .iter() .filter(|v| !v.name.is_empty() && !v.delete_marker) .map(|v| ObjectVersion { - key: Some(v.name.to_owned()), + key: Some(encode_output_value(&v.name)), last_modified: v.mod_time.map(Timestamp::from), size: Some(v.size), version_id: Some(v.version_id.map(|id| id.to_string()).unwrap_or_else(|| "null".to_string())), @@ -197,7 +216,7 @@ pub(crate) fn build_list_object_versions_output( .iter() .filter(|o| o.delete_marker) .map(|o| DeleteMarkerEntry { - key: Some(o.name.to_owned()), + key: Some(encode_output_value(&o.name)), version_id: Some(o.version_id.map(|id| id.to_string()).unwrap_or_else(|| "null".to_string())), is_latest: Some(o.is_latest), last_modified: o.mod_time.map(Timestamp::from), @@ -208,24 +227,32 @@ pub(crate) fn build_list_object_versions_output( let common_prefixes: Vec = object_infos .prefixes .into_iter() - .map(|v| CommonPrefix { prefix: Some(v) }) + .map(|v| CommonPrefix { + prefix: Some(encode_output_value(&v)), + }) .collect(); // Only return markers when they are non-empty to preserve S3 client compatibility. - let next_key_marker = object_infos.next_marker.filter(|v| !v.is_empty()); + let next_key_marker = object_infos + .next_marker + .filter(|v| !v.is_empty()) + .map(|marker| encode_output_value(&marker)); let next_version_id_marker = object_infos.next_version_idmarker.filter(|v| !v.is_empty()); ListObjectVersionsOutput { is_truncated: Some(object_infos.is_truncated), - max_keys: Some(max_keys), - delimiter, + max_keys: Some(params.max_keys), + delimiter: params.delimiter.as_deref().map(encode_output_value), + encoding_type: encoding_type.cloned(), + key_marker: Some(encode_output_value(params.key_marker.as_deref().unwrap_or_default())), name: Some(bucket), - prefix: Some(prefix), + prefix: Some(encode_output_value(¶ms.prefix)), common_prefixes: Some(common_prefixes), versions: Some(versions), delete_markers: Some(delete_markers), next_key_marker, next_version_id_marker, + version_id_marker: Some(params.version_id_marker.clone().unwrap_or_default()), ..Default::default() } } @@ -244,14 +271,7 @@ pub(crate) fn build_list_objects_v2_output( ) -> ListObjectsV2Output { // Apply URL encoding if encoding_type is "url". // S3 URL encoding encodes special characters but keeps '/' unencoded. - let should_encode = encoding_type.as_ref().is_some_and(|e| e.as_str() == EncodingType::URL); - - let encode_s3_name = |name: &str| -> String { - name.split('/') - .map(|part| encode(part).to_string()) - .collect::>() - .join("/") - }; + let should_encode = should_encode_url(encoding_type.as_ref()); let objects: Vec = object_infos .objects @@ -336,11 +356,19 @@ fn calculate_next_marker(v2: &ListObjectsV2Output) -> Option { .and_then(|prefix| prefix.prefix.as_ref()) .cloned(); - // NextMarker should be the lexicographically last item. - // This matches S3 standard behavior. + let sort_value = |value: &str| { + if should_encode_url(v2.encoding_type.as_ref()) { + percent_decode_str(value).decode_utf8_lossy().into_owned() + } else { + value.to_owned() + } + }; + + // NextMarker should be selected by raw object name ordering, even when the + // response fields are URL encoded. match (last_key, last_prefix) { (Some(k), Some(p)) => { - if k > p { + if sort_value(&k) > sort_value(&p) { Some(k) } else { Some(p) @@ -355,8 +383,8 @@ fn calculate_next_marker(v2: &ListObjectsV2Output) -> Option { #[cfg(test)] mod tests { use super::{ - build_list_buckets_output, build_list_object_versions_output, build_list_objects_output, build_list_objects_v2_output, - parse_list_object_versions_params, parse_list_objects_v2_params, + ListObjectVersionsParams, build_list_buckets_output, build_list_object_versions_output, build_list_objects_output, + build_list_objects_v2_output, parse_list_object_versions_params, parse_list_objects_v2_params, }; use crate::storage::s3_api::common::rustfs_owner; use rustfs_ecstore::store_api::{BucketInfo, ListObjectVersionsInfo, ListObjectsV2Info, ObjectInfo}; @@ -401,6 +429,19 @@ mod tests { assert_eq!(output.marker, Some("m-1".to_string())); } + #[test] + fn test_list_objects_marker_preserves_request_value_when_url_encoding_requested() { + let output = build_list_objects_output( + ListObjectsV2Output { + encoding_type: Some(EncodingType::from_static(EncodingType::URL)), + ..Default::default() + }, + Some("logs and more/start after".to_string()), + ); + + assert_eq!(output.marker.as_deref(), Some("logs and more/start after")); + } + #[test] fn test_list_objects_marker_defaults_to_empty_string() { let output = build_list_objects_output(ListObjectsV2Output::default(), None); @@ -425,6 +466,25 @@ mod tests { assert_eq!(output.next_marker, Some("zebra/".to_string())); } + #[test] + fn test_list_objects_next_marker_compares_raw_values_when_url_encoded() { + let v2 = ListObjectsV2Output { + is_truncated: Some(true), + encoding_type: Some(EncodingType::from_static(EncodingType::URL)), + contents: Some(vec![Object { + key: Some("z".to_string()), + ..Default::default() + }]), + common_prefixes: Some(vec![CommonPrefix { + prefix: Some("%7E/".to_string()), + }]), + ..Default::default() + }; + + let output = build_list_objects_output(v2, None); + assert_eq!(output.next_marker, Some("%7E/".to_string())); + } + #[test] fn test_list_objects_next_marker_is_none_when_not_truncated() { let v2 = ListObjectsV2Output { @@ -478,11 +538,11 @@ mod tests { true, 1000, "bucket-b".to_string(), - "prefix-b".to_string(), + "prefix b/child".to_string(), Some("/".to_string()), Some(EncodingType::from_static(EncodingType::URL)), None, - None, + Some("start after".to_string()), ); let contents = output.contents.as_ref().expect("contents should exist"); @@ -491,6 +551,9 @@ mod tests { assert_eq!(contents[0].key.as_deref(), Some("dir%20a/file%2Bb%25.txt")); assert_eq!(common_prefixes[0].prefix.as_deref(), Some("prefix%20a/sub%2B")); + assert_eq!(output.prefix.as_deref(), Some("prefix b/child")); + assert_eq!(output.delimiter.as_deref(), Some("/")); + assert_eq!(output.start_after.as_deref(), Some("start after")); assert!(contents[0].owner.is_some()); assert_eq!( contents[0].owner.as_ref().and_then(|owner| owner.display_name.clone()), @@ -515,11 +578,11 @@ mod tests { "prefix-c".to_string(), None, None, - Some(String::new()), + Some("echo-token".to_string()), Some("start-after".to_string()), ); - assert_eq!(output.continuation_token, Some(String::new())); + assert_eq!(output.continuation_token, Some("echo-token".to_string())); assert_eq!(output.start_after, Some("start-after".to_string())); assert_eq!( output.next_continuation_token, @@ -549,7 +612,7 @@ mod tests { #[test] fn test_parse_list_objects_v2_params_defaults_and_echo_behavior() { - let parsed = parse_list_objects_v2_params(None, Some(String::new()), None, Some(String::new()), Some(String::new())) + let parsed = parse_list_objects_v2_params(None, Some(String::new()), None, None, Some(String::new())) .expect("parse should succeed"); assert_eq!(parsed.prefix, String::new()); @@ -557,7 +620,7 @@ mod tests { assert_eq!(parsed.delimiter, None); assert_eq!(parsed.response_start_after, Some(String::new())); assert_eq!(parsed.start_after_for_query, None); - assert_eq!(parsed.response_continuation_token, Some(String::new())); + assert_eq!(parsed.response_continuation_token, None); assert_eq!(parsed.decoded_continuation_token, None); } @@ -583,6 +646,15 @@ mod tests { assert_eq!(*err.code(), S3ErrorCode::InvalidArgument); } + #[test] + fn test_parse_list_objects_v2_params_preserves_empty_continuation_token() { + let parsed = parse_list_objects_v2_params(None, None, Some(1000), Some(String::new()), None) + .expect("empty continuation token should be accepted"); + + assert_eq!(parsed.response_continuation_token, Some(String::new())); + assert_eq!(parsed.decoded_continuation_token, None); + } + #[test] fn test_parse_list_objects_v2_params_decodes_continuation_token() { let raw = "token-123"; @@ -646,9 +718,14 @@ mod tests { let output = build_list_object_versions_output( object_infos, "bucket-a".to_string(), - "prefix-a".to_string(), - Some("/".to_string()), - 123, + &ListObjectVersionsParams { + prefix: "prefix-a".to_string(), + delimiter: Some("/".to_string()), + key_marker: Some("marker-a".to_string()), + version_id_marker: Some("version-marker-a".to_string()), + max_keys: 123, + }, + None, ); assert_eq!(output.is_truncated, Some(true)); @@ -656,6 +733,8 @@ mod tests { assert_eq!(output.name, Some("bucket-a".to_string())); assert_eq!(output.prefix, Some("prefix-a".to_string())); assert_eq!(output.delimiter, Some("/".to_string())); + assert_eq!(output.key_marker, Some("marker-a".to_string())); + assert_eq!(output.version_id_marker, Some("version-marker-a".to_string())); assert_eq!(output.next_key_marker, None); assert_eq!(output.next_version_id_marker, Some("next-version-id".to_string())); @@ -678,6 +757,58 @@ mod tests { ); } + #[test] + fn test_build_list_object_versions_output_url_encodes_response_fields() { + let object_infos = ListObjectVersionsInfo { + is_truncated: true, + next_marker: Some("logs and more/next marker".to_string()), + next_version_idmarker: Some("next-version-id".to_string()), + objects: vec![ + ObjectInfo { + name: "logs and more/file name.txt".to_string(), + version_id: Some(Uuid::nil()), + ..Default::default() + }, + ObjectInfo { + name: "logs and more/delete marker.txt".to_string(), + delete_marker: true, + ..Default::default() + }, + ], + prefixes: vec!["logs and more/sub prefix/".to_string()], + }; + + let output = build_list_object_versions_output( + object_infos, + "bucket-a".to_string(), + &ListObjectVersionsParams { + prefix: "logs and more/".to_string(), + delimiter: Some("/".to_string()), + key_marker: Some("logs and more/start marker".to_string()), + version_id_marker: Some("version-marker-a".to_string()), + max_keys: 123, + }, + Some(&EncodingType::from_static(EncodingType::URL)), + ); + + assert_eq!(output.encoding_type.as_ref().map(EncodingType::as_str), Some(EncodingType::URL)); + assert_eq!(output.prefix.as_deref(), Some("logs%20and%20more/")); + assert_eq!(output.delimiter.as_deref(), Some("/")); + assert_eq!(output.key_marker.as_deref(), Some("logs%20and%20more/start%20marker")); + assert_eq!(output.next_key_marker.as_deref(), Some("logs%20and%20more/next%20marker")); + assert_eq!(output.version_id_marker.as_deref(), Some("version-marker-a")); + assert_eq!(output.next_version_id_marker.as_deref(), Some("next-version-id")); + + let versions = output.versions.unwrap_or_default(); + assert_eq!(versions[0].key.as_deref(), Some("logs%20and%20more/file%20name.txt")); + + let delete_markers = output.delete_markers.unwrap_or_default(); + assert_eq!(delete_markers[0].key.as_deref(), Some("logs%20and%20more/delete%20marker.txt")); + + let prefixes = output.common_prefixes.unwrap_or_default(); + assert_eq!(prefixes[0].prefix.as_deref(), Some("logs%20and%20more/sub%20prefix/")); + } + fn object_info(name: &str) -> ObjectInfo { ObjectInfo { name: name.to_string(),