fix(s3): preserve listing pagination parity (#3117)

* fix: preserve S3 listing pagination parity

S3 listing responses need stable wire semantics across encoded listings, version markers, and quorum-sensitive metadata reads. This tightens marker handling, response encoding, and version-marker pagination while keeping the changes scoped to listing paths and regression tests.

Constraint: Branch, code, commit, and PR text must avoid restricted upstream project naming.
Constraint: Verification required Rust 1.95 toolchain path because Homebrew cargo-clippy resolved to 1.94.
Rejected: Treat null version-id-marker as no marker | repeats or skips the marker boundary for null-version listings.
Rejected: Compare encoded next-marker values directly | encoded output can diverge from raw listing order.
Confidence: high
Scope-risk: moderate
Directive: Do not change listing marker semantics without covering V1, V2, versions, encoded responses, and null-version markers together.
Tested: cargo fmt --all --check
Tested: cargo clippy --workspace --all-features --all-targets -- -D warnings
Tested: make pre-commit with LC_ALL=en_US.UTF-8 and Rust 1.95 PATH
Tested: git diff --check
Not-tested: live distributed object-store compatibility test against a remote cluster

* fix: keep listing echo fields raw

S3 compatibility tests expect list response echo fields such as Prefix, Delimiter, StartAfter, and Marker to preserve the request value even when URL encoding is requested. Keep URL encoding scoped to object keys and common prefixes while preserving the next-marker raw comparison fix.

Constraint: PR branch was updated with latest main before this fix.

Constraint: Do not use restricted upstream project naming in commit text.

Rejected: Encode all response string fields | breaks compatibility tests for unreadable prefix values.

Confidence: high

Scope-risk: narrow

Tested: cargo test -p rustfs list_objects --lib

Tested: cargo fmt --all --check

Tested: cargo clippy --workspace --all-features --all-targets -- -D warnings

Tested: make pre-commit

Tested: git diff --check

Not-tested: live CI s3 compatibility rerun before push

* fix: accept empty listing continuation token

S3 compatibility tests treat an empty ListObjectsV2 continuation token as an explicit empty echo value, not as an invalid base64 token. Preserve the request echo while skipping decoded-token pagination for the empty string.

Constraint: Keep invalid non-empty continuation tokens rejected before store lookup.

Rejected: Drop empty continuation tokens from the response | compatibility tests assert the empty echo field is present.

Confidence: high

Scope-risk: narrow

Tested: cargo test -p rustfs list_objects --lib

Tested: cargo fmt --all --check

Tested: cargo clippy --workspace --all-features --all-targets -- -D warnings

Tested: make pre-commit

Tested: git diff --check

Not-tested: live CI s3 compatibility rerun after this push

* fix: delete explicit null object versions

S3 compatibility cleanup lists unversioned objects with VersionId=null and then sends that value back through DeleteObjects. The API layer keeps null as an internal sentinel, but the storage delete path must map it back to the stored null version instead of treating it as a real UUID.

Constraint: Keep the wire response able to echo null version IDs.

Rejected: Treat VersionId=null the same as an absent version id everywhere | versioned and suspended buckets need explicit null version semantics.

Confidence: high

Scope-risk: moderate

Tested: cargo test -p rustfs-ecstore delete_file_info_version_id_maps_explicit_null_version_to_stored_null

Tested: cargo test -p rustfs normalize_delete_objects_version_id_preserves_explicit_null_marker --lib

Tested: cargo fmt --all --check

Tested: cargo clippy --workspace --all-features --all-targets -- -D warnings

Tested: make pre-commit

Tested: git diff --check

Not-tested: live CI s3 compatibility rerun after this push

* fix: keep version marker scoped to marker key

Version listing cleanup can delete the key returned as the previous page marker before requesting the next page. When the marker key is no longer present, the version marker must not be applied to the first later key, otherwise each cleanup page can skip one null-version object and leave the bucket non-empty.

Constraint: Preserve version-marker pagination for marker keys that still exist.

Rejected: Drop version markers whenever the marker key was supplied | multi-version marker keys still need intra-key pagination.

Confidence: high

Scope-risk: narrow

Tested: cargo test -p rustfs-ecstore version_marker_is_applied_only_when_key_marker_entry_is_present

Tested: cargo test -p rustfs-ecstore delete_file_info_version_id_maps_explicit_null_version_to_stored_null

Tested: cargo test -p rustfs normalize_delete_objects_version_id_preserves_explicit_null_marker --lib

Tested: cargo fmt --all --check

Tested: cargo clippy --workspace --all-features --all-targets -- -D warnings

Tested: make pre-commit

Tested: git diff --check

Not-tested: live CI s3 compatibility rerun after this push

---------

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: loverustfs <hello@rustfs.com>
This commit is contained in:
weisd
2026-05-30 19:19:59 +08:00
committed by GitHub
parent 8e809f005d
commit ca4793f93e
6 changed files with 686 additions and 152 deletions
+28 -2
View File
@@ -588,6 +588,18 @@ impl SetDisks {
// shuffle_disks TODO: use origin value
}
fn is_explicit_null_version(version_id: Option<Uuid>) -> bool {
version_id == Some(Uuid::nil())
}
fn delete_file_info_version_id(version_id: Option<Uuid>) -> Option<Uuid> {
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
+115 -7
View File
@@ -615,11 +615,12 @@ impl ObjectInfo {
bucket: &str,
prefix: &str,
delimiter: Option<String>,
after_version_id: Option<Uuid>,
after_version_marker: Option<VersionMarker>,
) -> Vec<ObjectInfo> {
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 {
+355 -99
View File
@@ -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<String>,
// Marker to resume listing.
// The response will be the first entry >= this object name.
pub marker: Option<String>,
// 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<VersionMarker> {
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<VersionMarker>,
) -> Option<VersionMarker> {
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::<usize>() {
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::<usize>() {
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::<usize>() {
Ok(res) => self.pool_idx = Some(res),
Err(_) => {
self.id = Some(Uuid::new_v4().to_string());
self.create = true;
}
},
"s" => match value.parse::<usize>() {
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::<Vec<_>>();
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::<Vec<_>>();
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<Uuid> = 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<Uuid> = 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<Uuid> = 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<Uuid> = 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)])
+11 -9
View File
@@ -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,
&params.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, &params, encoding_type.as_ref());
Ok(S3Response::new(output))
}
+11
View File
@@ -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 {
+166 -35
View File
@@ -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::<Vec<_>>()
.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<String>,
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<ObjectVersion> = 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<CommonPrefix> = 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(&params.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::<Vec<_>>()
.join("/")
};
let should_encode = should_encode_url(encoding_type.as_ref());
let objects: Vec<Object> = object_infos
.objects
@@ -336,11 +356,19 @@ fn calculate_next_marker(v2: &ListObjectsV2Output) -> Option<String> {
.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<String> {
#[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(),