diff --git a/docs/operations/on-demand-migration.md b/docs/operations/on-demand-migration.md index 06d066fd4..137f42248 100644 --- a/docs/operations/on-demand-migration.md +++ b/docs/operations/on-demand-migration.md @@ -28,6 +28,8 @@ The two rollout switches have different defaults. Unset or invalid boolean value - `RUSTFS_ON_DEMAND_MIGRATION_LIST_V2_TOKENS` defaults to `false` and allows a v1 listing to first issue a v2 token after an empty truncated merged page. Existing v2 tokens keep their budget even on reader-only nodes. Consuming an object/common prefix or reaching a new EOF resets the budget to v1 without changing the chain's framing. - `RUSTFS_ON_DEMAND_MIGRATION_LIST_FRAMED_TOKENS` defaults to `true`, preserving the framed output of the #7187 / `e1608fbd9` generation. It allows a bare/new merged listing to first issue a NUL-prefixed JSON envelope inside the existing base64 encoding. Existing framed chains stay framed even with this switch off, including a reset to v1 and local continuation after list-through is disabled. With framing issuance off, new bare v1 output keeps its historical bytes; ordinary local listings remain unchanged. This switch does not enable the v2 budget. +When list-through or the global module is disabled, existing merged continuations retain their last emitted key and original bare/framed format, filter already returned objects and common prefixes, and make no source requests. If the local-only response remains truncated, its continuation can resume both sides after re-enabling list-through. + Choose the upgrade path from the binaries currently serving LIST requests, including every load-balancer route. This build reads complete historical bare envelopes and framed v1/v2 envelopes with the same strict version/count validation. There is no single writer format understood by both bare-only and framed-only readers. A bare-v1-only binary rejects bare v2 with `400 InvalidArgument`; a bare-only reader mistakes framed input for a local marker, while a framed-only reader mistakes bare input for one. These framing mismatches can restart a merged scan and lose its budget without returning an error. For example, with local keys `b,d` and source keys `a,c`, a new node issuing bare after `a` followed by an `e1608fbd9` reader can return `a` again. That old reader then emits framing, so the symptom need not be an infinite loop. - **Upgrading from the framed-only #7187 / `e1608fbd9` generation:** before deploying the first new node, set `RUSTFS_ON_DEMAND_MIGRATION_LIST_FRAMED_TOKENS=true` in the new nodes' deployment environment, or leave it unset to use this build's `true` default. Remove any previous explicit `false` override. Existing framed-only nodes ignore this new variable and already emit framing. Keep framing enabled while any of those nodes serves continuations. While list-through remains active, new and existing framed v1/v2 chains then retain their cursors and any active budget in both directions. Only after all serving nodes are dual readers may you choose `false`; existing framed chains still stay framed, while newly started bare chains must stay on dual readers. diff --git a/rustfs/src/app/bucket_list_through.rs b/rustfs/src/app/bucket_list_through.rs index 7fecdfce1..40713b952 100644 --- a/rustfs/src/app/bucket_list_through.rs +++ b/rustfs/src/app/bucket_list_through.rs @@ -33,10 +33,9 @@ use super::storage_api::bucket_usecase::s3_api::bucket::ListObjectsV2Params; use crate::app::object::shared::{odm_source_unavailable_error, odm_state_error_class}; use crate::error::ApiError; use crate::on_demand_migration::{ - BucketOdmState, LIST_THROUGH_TOKEN_VERSION, ListEntryKey, ListPageError, ListThroughCursor, ListThroughMerger, - ListThroughToken, ListThroughTokenError, MergeSide, OnDemandMigrationSys, SOURCE_LIST_MAX_RATE_WAIT, SourceClient, - SourceError, SourceErrorPolicy, SourceListPlan, SourceListRequest, SourceObject, SourcePage, decode_continuation_token, - source_list_plan, + BucketOdmState, ListEntryKey, ListPageError, ListThroughCursor, ListThroughMerger, ListThroughToken, ListThroughTokenError, + MergeSide, OnDemandMigrationSys, SOURCE_LIST_MAX_RATE_WAIT, SourceClient, SourceError, SourceErrorPolicy, SourceListPlan, + SourceListRequest, SourceObject, SourcePage, decode_continuation_token, source_list_plan, }; use futures::StreamExt; use rustfs_utils::http::{SUFFIX_SOURCE_PROXY_REQUEST, get_header}; @@ -61,15 +60,6 @@ const ENV_LIST_FRAMED_TOKENS: &str = "RUSTFS_ON_DEMAND_MIGRATION_LIST_FRAMED_TOK /// source-only keys for a shadowing delete marker. const DELETE_MARKER_PROBE_CONCURRENCY: usize = 32; -/// Where the local side of a listing resumes. -pub(crate) enum LocalListCursor { - /// Continuation token for the local store, `None` for the first page. - Token(Option), - /// The local side of a merged listing was already exhausted, so a listing - /// that no longer merges has nothing left to return. - Exhausted, -} - /// Reads the (already base64-decoded) continuation token. /// /// This runs whether or not the bucket merges: a token handed out while @@ -83,47 +73,6 @@ pub(crate) fn decode_list_cursor(decoded: Option<&str>) -> S3Result, merged: Option<&ListThroughToken>) -> LocalListCursor { - match merged { - Some(token) if token.local_done => LocalListCursor::Exhausted, - Some(token) => LocalListCursor::Token(token.local.clone()), - None => LocalListCursor::Token(decoded.map(str::to_string)), - } -} - -/// A framed chain keeps its envelope when the bucket stops consulting source. -pub(crate) fn preserve_framed_local_cursor(info: &mut ListObjectsV2Info, previous: Option<&ListThroughToken>) { - let Some(previous) = previous.filter(|token| token.framed) else { - return; - }; - let Some(next) = info.next_continuation_token.take() else { - return; - }; - let mut token = previous.clone(); - token.local = Some(next); - token.local_done = false; - if let Some(last_key) = info - .objects - .iter() - .map(|object| object.name.as_str()) - .chain(info.prefixes.iter().map(String::as_str)) - .max() - { - token.last_key = Some( - token - .last_key - .as_deref() - .map_or(last_key, |previous| previous.max(last_key)) - .to_string(), - ); - token.v = LIST_THROUGH_TOKEN_VERSION; - token.no_progress = None; - } - info.next_continuation_token = Some(token.encode()); -} - fn invalid_continuation_token(err: &ListThroughTokenError) -> S3Error { debug!(error = %err, "rejected an on-demand migration list continuation token"); S3Error::with_message(S3ErrorCode::InvalidArgument, "Invalid continuation token".to_string()) @@ -240,7 +189,10 @@ fn source_page_entries(bucket: &str, page: SourcePage) -> Vec { interleave(objects, page.common_prefixes) } -/// Runs one merged `ListObjectsV2` page. +/// Runs one merged `ListObjectsV2` page. With no source state, continues only +/// the local side of an existing merged chain. The original local page cursor +/// is reread and raw names are filtered against the last emitted key here; +/// sending that key to storage as a marker could interpret literal cache tags. /// /// Cost: at most two listings per side per request — the first page of each /// side, plus one refill when the previous page had already consumed most of @@ -249,30 +201,34 @@ fn source_page_entries(bucket: &str, page: SourcePage) -> Vec { /// source listings. pub(crate) async fn merged_list_objects_v2( store: &Arc, - state: &Arc, + state: Option<&Arc>, bucket: &str, params: &ListObjectsV2Params, fetch_owner: bool, incl_deleted: bool, token: Option<&ListThroughToken>, ) -> S3Result { - let policy = &state.config().policy; let max_keys = usize::try_from(params.max_keys).unwrap_or(0); let mut merger = ListThroughMerger::new(max_keys, token); let mut buffers: [Vec>; 2] = [Vec::new(), Vec::new()]; let mut degraded = false; - let plan = source_list_plan(¶ms.prefix, state.config().filter.prefix.as_deref(), params.delimiter.as_deref()); + let plan = state.map_or(SourceListPlan::Skip, |state| { + source_list_plan(¶ms.prefix, state.config().filter.prefix.as_deref(), params.delimiter.as_deref()) + }); let mut client: Option> = None; - match &plan { - // Nothing the source holds can appear under this prefix; that is a - // filter decision, not a degradation. - SourceListPlan::Skip => merger.disable_source(), - _ => match state.client() { - Ok(ready) if state.breaker().allow_request() => client = Some(Arc::clone(ready)), - Ok(_) => degrade_or_fail(&mut merger, &mut degraded, policy.source_error, "breaker_open")?, - Err(error) => degrade_or_fail(&mut merger, &mut degraded, policy.source_error, odm_state_error_class(error))?, - }, + match (state, &plan) { + // A disabled or excluded source is not a degradation. The merger + // retains its unconsumed source cursor for a later enabled request. + (None, _) | (_, SourceListPlan::Skip) => merger.disable_source(), + (Some(state), _) => { + let policy = &state.config().policy; + match state.client() { + Ok(ready) if state.breaker().allow_request() => client = Some(Arc::clone(ready)), + Ok(_) => degrade_or_fail(&mut merger, &mut degraded, policy.source_error, "breaker_open")?, + Err(error) => degrade_or_fail(&mut merger, &mut degraded, policy.source_error, odm_state_error_class(error))?, + } + } } while let Some(fetch) = merger.next_fetch() { @@ -295,6 +251,8 @@ pub(crate) async fn merged_list_objects_v2( (interleave(objects, info.prefixes), info.is_truncated, info.next_continuation_token) } MergeSide::Source => { + let state = state.expect("the source side is disabled without bucket state"); + let policy = &state.config().policy; let client = client.as_ref().expect("the source side is disabled without a client"); match fetch_source_page(state, client, bucket, params, &plan, fetch.token.as_deref()).await { Ok(page) => page, @@ -314,6 +272,7 @@ pub(crate) async fn merged_list_objects_v2( if let Err(error) = merger.push_page(fetch.side, keys, is_truncated, next_token) { match fetch.side { MergeSide::Source => { + let policy = &state.expect("only an enabled source can return a page").config().policy; degrade_or_fail(&mut merger, &mut degraded, policy.source_error, "invalid_pagination")?; continue; } @@ -324,10 +283,15 @@ pub(crate) async fn merged_list_objects_v2( } let issue_progress_tokens = rustfs_utils::get_env_bool(ENV_LIST_PROGRESS_TOKENS, false); - let framed = token.is_some_and(|token| token.framed) || rustfs_utils::get_env_bool(ENV_LIST_FRAMED_TOKENS, true); + let framed = + token.is_some_and(|token| token.framed) || (state.is_some() && rustfs_utils::get_env_bool(ENV_LIST_FRAMED_TOKENS, true)); let outcome = match merger.finish(issue_progress_tokens) { Ok(outcome) => outcome, Err(ListPageError::NoProgress(MergeSide::Source)) => { + let policy = &state + .expect("a disabled source cannot exhaust the progress budget") + .config() + .policy; degrade_or_fail(&mut merger, &mut degraded, policy.source_error, "invalid_pagination")?; merger .finish(issue_progress_tokens) @@ -353,7 +317,10 @@ pub(crate) async fn merged_list_objects_v2( } } - if !source_only_keys.is_empty() && policy.respect_local_delete_marker && bucket_keeps_delete_markers(bucket).await { + if !source_only_keys.is_empty() + && state.is_some_and(|state| state.config().policy.respect_local_delete_marker) + && bucket_keeps_delete_markers(bucket).await + { let shadowed = local_delete_markers(store, bucket, &source_only_keys).await; objects.retain(|object| !shadowed.contains(&object.name)); } @@ -586,14 +553,18 @@ mod tests { let encoded = resume.encode(); let decoded = decode_list_cursor(Some(&encoded)).expect("a valid envelope decodes"); assert_eq!(decoded.as_ref(), Some(&resume)); - assert!(matches!( - local_cursor(Some(&encoded), decoded.as_ref()), - LocalListCursor::Token(Some(local)) if local == "local-2" - )); + let mut merger = ListThroughMerger::new(2, decoded.as_ref()); + merger.disable_source(); + let fetch = merger.next_fetch().expect("unconsumed local page"); + assert_eq!(fetch.side, MergeSide::Local); + assert_eq!(fetch.token.as_deref(), Some("local-2")); let encoded = token(None, true).encode(); let decoded = decode_list_cursor(Some(&encoded)).expect("a valid envelope decodes"); - assert!(matches!(local_cursor(Some(&encoded), decoded.as_ref()), LocalListCursor::Exhausted)); + let mut merger = ListThroughMerger::new(2, decoded.as_ref()); + merger.disable_source(); + assert!(merger.next_fetch().is_none()); + assert!(!merger.finish(false).expect("local EOF").is_truncated); } #[test] @@ -606,31 +577,44 @@ mod tests { let encoded = resume.encode(); let decoded = decode_list_cursor(Some(&encoded)).expect("a v2 envelope decodes"); assert_eq!(decoded.as_ref(), Some(&resume)); - assert!(matches!( - local_cursor(Some(&encoded), decoded.as_ref()), - LocalListCursor::Token(Some(local)) if local == "local-2" - )); + let mut merger = ListThroughMerger::new(2, decoded.as_ref()); + merger.disable_source(); + let fetch = merger.next_fetch().expect("unconsumed local page"); + assert_eq!(fetch.side, MergeSide::Local); + assert_eq!(fetch.token.as_deref(), Some("local-2")); resume.local_done = true; let encoded = resume.encode(); let decoded = decode_list_cursor(Some(&encoded)).expect("v2 with local EOF decodes"); - assert!(matches!(local_cursor(Some(&encoded), decoded.as_ref()), LocalListCursor::Exhausted)); + let mut merger = ListThroughMerger::new(2, decoded.as_ref()); + merger.disable_source(); + assert!(merger.next_fetch().is_none()); + assert!(!merger.finish(false).expect("local EOF").is_truncated); } } #[test] fn a_plain_local_token_is_passed_through_and_a_tampered_one_is_rejected() { + use crate::app::storage_api::bucket_usecase::s3_api::bucket::parse_list_objects_v2_params; + let json_key = r#"{"t":"odm-list","v":1,"local_done":true}"#; assert!(decode_list_cursor(Some(json_key)).expect("valid local key").is_none()); - assert!(matches!(local_cursor(Some(json_key), None), LocalListCursor::Token(Some(local)) if local == json_key)); assert!( decode_list_cursor(Some("photos/a.jpg")) .expect("plain markers decode") .is_none() ); - assert!(matches!( - local_cursor(Some("photos/a.jpg"), None), - LocalListCursor::Token(Some(local)) if local == "photos/a.jpg" - )); + + for marker in [json_key, "photos/a.jpg"] { + let params = parse_list_objects_v2_params( + None, + None, + Some(2), + Some(base64_simd::STANDARD.encode_to_string(marker.as_bytes())), + None, + ) + .expect("plain local listing parameters"); + assert_eq!(params.decoded_continuation_token.as_deref(), Some(marker)); + } let tampered = token(Some("local-2"), false).encode().replace("\"v\":1", "\"v\":9"); let err = decode_list_cursor(Some(&tampered)).expect_err("a bumped version is rejected"); @@ -1783,7 +1767,7 @@ mod tests { Duration::from_secs(10), merged_list_objects_v2( &store, - &state, + Some(&state), &input.bucket, ¶ms, input.fetch_owner.unwrap_or_default(), @@ -2116,71 +2100,303 @@ mod tests { .expect("merged continuation token") } + #[derive(Clone, Copy)] + enum DisabledListScenario { + FirstLocalPage, + PartiallyConsumedLocalPage, + CacheTag(&'static str), + CommonPrefixes, + Reenable, + } + + async fn assert_disabled_list_progress(scenario: DisabledListScenario, framed: bool, disable_module: bool) { + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; + + let (local_keys, source_names, expected_first, expected_disabled) = match scenario { + DisabledListScenario::FirstLocalPage => (vec!["a", "c"], ["b", "d"], ["a", "b"], vec!["c", "z-local"]), + DisabledListScenario::PartiallyConsumedLocalPage => { + (vec!["a", "b", "c", "e"], ["d", "f"], ["a", "b"], vec!["e", "z-local"]) + } + DisabledListScenario::CacheTag(key) => (vec!["a", "b!", "b0", "c"], [key, "f"], ["a", "b!"], vec!["c", "z-local"]), + DisabledListScenario::CommonPrefixes => (vec!["p/a/1", "p/c/1"], ["p/b/", "p/d/"], ["p/a/", "p/b/"], vec!["p/c/"]), + DisabledListScenario::Reenable => (vec!["a", "c", "e"], ["b", "f"], ["a", "b"], vec!["c", "e"]), + }; + let prefixes = matches!(scenario, DisabledListScenario::CommonPrefixes); + let entries = source_names + .iter() + .map(|name| { + if prefixes { + format!("{name}") + } else { + format!("{name}1") + } + }) + .collect::(); + let body = format!( + "false{entries}" + ); + let source_disabled = Arc::new(AtomicBool::new(false)); + let source_requests = Arc::new(AtomicUsize::new(0)); + let forbidden = Arc::clone(&source_disabled); + let count = Arc::clone(&source_requests); + let (endpoint, server, stop) = list_source_with_response(std::iter::repeat(body), move |_, body| { + assert!(!forbidden.load(Ordering::SeqCst), "disabled listing must not call the source"); + count.fetch_add(1, Ordering::SeqCst); + body + }) + .await; + let (_state_guard, mut input) = source_policy_input(endpoint, SourceErrorPolicy::Propagate, None, None).await; + let store = shared_gating_ecstore().await; + for key in local_keys { + store + .put_object( + &input.bucket, + key, + &mut StoragePutObjReader::from_vec(vec![1]), + &StorageObjectOptions::default(), + ) + .await + .expect("seed disabled-list local keys"); + } + if prefixes { + input.prefix = Some("p/".to_string()); + input.delimiter = Some("/".to_string()); + } + input.start_after = Some(if prefixes { "p/0" } else { "0" }.to_string()); + let page_names = |page: &ListObjectsV2Output| { + let mut names = page + .contents + .iter() + .flatten() + .map(|object| object.key.clone().expect("listed object key")) + .chain( + page.common_prefixes + .iter() + .flatten() + .map(|prefix| prefix.prefix.clone().expect("listed common prefix")), + ) + .collect::>(); + names.sort(); + names + }; + let first = execute_source_list(input.clone()).await.expect("first merged page").output; + assert_eq!(page_names(&first), expected_first); + assert_eq!(first.key_count, Some(2)); + assert_eq!(first.is_truncated, Some(true)); + let mut cursor = first.next_continuation_token.expect("merged cursor"); + let mut token = decode_wire_token(&cursor); + assert_eq!(token.framed, framed); + let second_page = match scenario { + DisabledListScenario::PartiallyConsumedLocalPage => Some((["c", "d"], ["c", "e"])), + DisabledListScenario::CacheTag(key) => Some((["b0", key], ["b0", "c"])), + _ => None, + }; + if let Some((expected_second, expected_replay)) = second_page { + let first_local = token.local.clone().expect("first local page was fully consumed"); + assert!( + first_local.starts_with(&format!("{}[rustfs_cache:", expected_first[1])), + "expected a real opaque cache cursor: {first_local}" + ); + input.continuation_token = Some(cursor); + let second = execute_source_list(input.clone()).await.expect("second merged page").output; + assert_eq!(page_names(&second), expected_second); + cursor = second.next_continuation_token.expect("partially consumed local page"); + token = decode_wire_token(&cursor); + assert_eq!( + token.local.as_deref(), + Some(first_local.as_str()), + "keep the partially consumed local page" + ); + assert_eq!(token.last_key.as_deref(), Some(expected_second[1])); + // ECStore prioritizes its opaque continuation over StartAfter. A + // fix that merely passes both would still replay c from this page. + let token_wins = Arc::clone(&store) + .list_objects_v2( + &input.bucket, + "", + Some(first_local), + None, + 2, + false, + Some(expected_second[1].to_string()), + false, + ) + .await + .expect("verify the real storage continuation contract"); + assert_eq!( + token_wins + .objects + .iter() + .map(|object| object.name.as_str()) + .collect::>(), + expected_replay + ); + } else { + assert_eq!(token.local, None, "the first local page remains partially consumed"); + assert_eq!(token.last_key.as_deref(), Some(expected_first[1])); + } + assert!(!token.local_done); + assert!(!token.source_done); + let requests_before_disable = source_requests.load(Ordering::SeqCst); + assert_eq!(requests_before_disable, if second_page.is_some() { 2 } else { 1 }); + let sys = OnDemandMigrationSys::get(); + let installed = sys.state(&input.bucket).expect("installed source state"); + let saved_config = installed.config().clone(); + source_disabled.store(true, Ordering::SeqCst); + if disable_module { + sys.set_module_enabled(false); + } else { + let mut disabled_config = saved_config.clone(); + disabled_config.policy.list_through = false; + sys.apply_for_incarnation(&input.bucket, installed.incarnation_id(), Some(&disabled_config)) + .await; + } + input.continuation_token = Some(cursor); + let disabled = execute_source_list(input.clone()).await.expect("local-only continuation"); + assert!(!disabled.headers.contains_key("x-rustfs-on-demand-migration-list")); + assert_eq!( + page_names(&disabled.output), + expected_disabled, + "framed={framed}, module={disable_module}" + ); + assert_eq!( + disabled.output.key_count, + Some(i32::try_from(expected_disabled.len()).expect("page size")) + ); + assert_eq!(disabled.output.prefix.as_deref(), Some(if prefixes { "p/" } else { "" })); + assert_eq!(disabled.output.delimiter, input.delimiter); + assert_eq!(disabled.output.start_after, input.start_after, "echo the client's original StartAfter"); + assert_eq!(source_requests.load(Ordering::SeqCst), requests_before_disable); + if matches!(scenario, DisabledListScenario::Reenable) { + assert_eq!(disabled.output.is_truncated, Some(true)); + let next = disabled.output.next_continuation_token.expect("local side still has z-local"); + let next_token = decode_wire_token(&next); + assert_eq!(next_token.framed, framed, "retain the chain's existing wire format"); + assert_eq!(next_token.source, token.source, "disabled requests do not consume source pages"); + assert_eq!(next_token.last_key.as_deref(), Some("e")); + source_disabled.store(false, Ordering::SeqCst); + if disable_module { + sys.set_module_enabled(true); + } else { + sys.apply_for_incarnation(&input.bucket, installed.incarnation_id(), Some(&saved_config)) + .await; + } + input.continuation_token = Some(next); + let resumed = execute_source_list(input).await.expect("reenabled continuation").output; + assert_eq!(page_names(&resumed), ["f", "z-local"], "both sides continue beyond the local-only page"); + assert_eq!(resumed.is_truncated, Some(false)); + assert!(resumed.next_continuation_token.is_none()); + assert_eq!(source_requests.load(Ordering::SeqCst), requests_before_disable + 1); + } else { + assert_eq!(disabled.output.is_truncated, Some(false)); + assert!(disabled.output.next_continuation_token.is_none()); + } + stop.cancel(); + let requests = tokio::time::timeout(Duration::from_secs(5), server) + .await + .expect("source server must finish") + .expect("source access must respect the disabled phase"); + assert_eq!(requests.len(), source_requests.load(Ordering::SeqCst)); + for request in requests { + assert!( + !request.contains("continuation-token="), + "the partially consumed source page is reread: {request}" + ); + if prefixes { + assert!(request.contains("prefix=p%2F") && request.contains("delimiter=%2F"), "{request}"); + } + } + } + + async fn disabled_list_progress_matrix(scenario: DisabledListScenario) { + for framed in [false, true] { + temp_env::async_with_vars( + [ + (ENV_LIST_PROGRESS_TOKENS, Some("true")), + (ENV_LIST_FRAMED_TOKENS, Some(if framed { "true" } else { "false" })), + ("RUSTFS_REPLICATION_ALLOW_LOOPBACK_TARGET", Some("true")), + ("HTTP_PROXY", None), + ("HTTPS_PROXY", None), + ("ALL_PROXY", None), + ("http_proxy", None), + ("https_proxy", None), + ("all_proxy", None), + ("NO_PROXY", Some("*")), + ("no_proxy", Some("*")), + ], + Box::pin(async move { + for disable_module in [false, true] { + assert_disabled_list_progress(scenario, framed, disable_module).await; + } + }), + ) + .await; + } + } + #[test] - fn framed_local_continuations_preserve_json_markers_and_zero_sized_budgets() { + #[serial_test::serial] + fn list_through_disabled_resumes_partially_consumed_first_local_page() { + run_large_stack_test("list-disabled-first", || { + disabled_list_progress_matrix(DisabledListScenario::FirstLocalPage) + }); + } + + #[test] + #[serial_test::serial] + fn list_through_disabled_resumes_partially_consumed_opaque_local_page() { + run_large_stack_test("list-disabled-opaque", || { + disabled_list_progress_matrix(DisabledListScenario::PartiallyConsumedLocalPage) + }); + } + + #[test] + #[serial_test::serial] + fn list_through_disabled_treats_source_cache_tag_names_as_literal_keys() { + run_large_stack_test("list-disabled-literal-cache-tag", || async { + for key in ["b[rustfs_cache:v1,return:]", "b[rustfs_cache:v2,return:]"] { + disabled_list_progress_matrix(DisabledListScenario::CacheTag(key)).await; + } + }); + } + + #[test] + #[serial_test::serial] + fn list_through_disabled_resumes_common_prefixes_without_repeating() { + run_large_stack_test("list-disabled-prefixes", || { + disabled_list_progress_matrix(DisabledListScenario::CommonPrefixes) + }); + } + + #[test] + #[serial_test::serial] + fn list_through_disabled_then_reenabled_retains_progress_on_both_sides() { + run_large_stack_test("list-disabled-reenable", || disabled_list_progress_matrix(DisabledListScenario::Reenable)); + } + + #[test] + fn local_merger_preserves_json_markers_and_resets_budget_on_progress() { let json_key = r#"{"t":"odm-list","v":1,"local_done":true}"#; - let mut resume = token(Some("local-2"), false); - resume.framed = true; - resume.v = 2; - resume.no_progress = Some(15); - let mut page = ListObjectsV2Info { - is_truncated: true, - next_continuation_token: Some(json_key.to_string()), - objects: vec![info(json_key)], - ..Default::default() - }; - preserve_framed_local_cursor(&mut page, Some(&resume)); - let raw = page.next_continuation_token.expect("local continuation"); - assert!(raw.starts_with("\0odm-list:")); - let decoded = decode_list_cursor(Some(&raw)) - .expect("framed local continuation") - .expect("envelope"); - assert!(decoded.framed); - assert_eq!( - decoded.local.as_deref(), - Some(json_key), - "the local marker is embedded without another encoding" - ); - assert_eq!(decoded.source, resume.source); - assert_eq!(decoded.last_key.as_deref(), Some(json_key)); - assert_eq!(decoded.v, 1); - assert_eq!(decoded.no_progress, None); - assert!(matches!(local_cursor(Some(&raw), Some(&decoded)), LocalListCursor::Token(Some(local)) if local == json_key)); - - let mut prefix_page = ListObjectsV2Info { - is_truncated: true, - next_continuation_token: Some("photos/".to_string()), - prefixes: vec!["photos/".to_string()], - ..Default::default() - }; - preserve_framed_local_cursor(&mut prefix_page, Some(&resume)); - let prefix = decode_list_cursor(prefix_page.next_continuation_token.as_deref()) - .expect("prefix continuation") - .expect("framed prefix envelope"); - assert!(prefix.framed); - assert_eq!(prefix.last_key.as_deref(), Some("photos/")); - assert_eq!(prefix.v, 1); - assert_eq!(prefix.no_progress, None); - - let mut zero = ListObjectsV2Info { - is_truncated: true, - next_continuation_token: resume.local.clone(), - ..Default::default() - }; - preserve_framed_local_cursor(&mut zero, Some(&resume)); - assert_eq!( - decode_list_cursor(zero.next_continuation_token.as_deref()).expect("zero-sized continuation"), - Some(resume.clone()) - ); - - resume.framed = false; - let mut ordinary = ListObjectsV2Info { - is_truncated: true, - next_continuation_token: Some(json_key.to_string()), - ..Default::default() - }; - preserve_framed_local_cursor(&mut ordinary, Some(&resume)); - assert_eq!(ordinary.next_continuation_token.as_deref(), Some(json_key)); + for key in [ListEntryKey::object(json_key), ListEntryKey::prefix("photos/")] { + let mut resume = token(Some("local-2"), false); + resume.v = 2; + resume.no_progress = Some(15); + let mut merger = ListThroughMerger::new(1, Some(&resume)); + merger.disable_source(); + merger + .push_page(MergeSide::Local, vec![key.clone()], true, Some(key.name.clone())) + .expect("valid local page"); + let outcome = merger.finish(false).expect("local progress resets the budget"); + assert_eq!(outcome.picks.len(), 1); + assert_eq!(outcome.picks[0].side, MergeSide::Local); + let next = outcome.next_token.expect("local continuation"); + assert_eq!(next.local.as_deref(), Some(key.name.as_str()), "the local marker is embedded unchanged"); + assert_eq!(next.source, resume.source); + assert_eq!(next.source_done, resume.source_done); + assert_eq!(next.last_key.as_deref(), Some(key.name.as_str())); + assert_eq!(next.v, 1); + assert_eq!(next.no_progress, None); + } } #[test] diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index f06ecf4b9..1cb5a8be9 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -2797,42 +2797,35 @@ impl DefaultBucketUsecase { false, ) } - (Some(state), _) => { + (None, None) => { + let infos = store + .list_objects_v2( + &bucket, + ¶ms.prefix, + params.decoded_continuation_token.clone(), + params.delimiter.clone(), + params.max_keys, + fetch_owner.unwrap_or_default(), + params.start_after_for_query.clone(), + incl_deleted, + ) + .await + .map_err(ApiError::from)?; + (infos, false) + } + (state, token) => { let outcome = list_through::merged_list_objects_v2( &store, - &state, + state.as_ref(), &bucket, ¶ms, fetch_owner.unwrap_or_default(), incl_deleted, - merged_token.as_ref(), + token, ) .await?; (outcome.info, outcome.degraded) } - (None, _) => { - let cursor = list_through::local_cursor(params.decoded_continuation_token.as_deref(), merged_token.as_ref()); - match cursor { - list_through::LocalListCursor::Exhausted => (StorageListObjectsV2Info::default(), false), - list_through::LocalListCursor::Token(token) => { - let mut infos = store - .list_objects_v2( - &bucket, - ¶ms.prefix, - token, - params.delimiter.clone(), - params.max_keys, - fetch_owner.unwrap_or_default(), - params.start_after_for_query.clone(), - incl_deleted, - ) - .await - .map_err(ApiError::from)?; - list_through::preserve_framed_local_cursor(&mut infos, merged_token.as_ref()); - (infos, false) - } - } - } }; let output = build_list_objects_v2_output(