From c5c2dc5aa0478a1ea7313a7d9cd70791e2a856d7 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sun, 26 Jul 2026 16:09:33 +0800 Subject: [PATCH] refactor(storage): resolve staged dead code in store module (#730) Drop the store module's blanket #![allow(dead_code)] and resolve every unit it was masking, instead of keeping them staged: - write_persistent_key_only_index: the keys-only writer wrapper only had test callers; the in-process rebuild flow already exists and uses the _with_metadata variant (prepare -> rebuild, triggered lazily from the opt-in listing path). Tests call _with_metadata directly now. - ListIndexLifecycle::recover_after_restart / mark_corrupt: production derives index health per request from the persisted artifacts (index file + namespace mutation journal), and restart recovery already lives in load_persistent_key_only_index's journal restore. The persisted- lifecycle design these transitions served was superseded, and Corrupt had no detector, so the never-constructed Corrupt state/reason variants go with them. - record_list_objects_index_opt_in_fallback (+ helpers): superseded by the inline per-site fallback recording in the opt-in listing path; the recorder recomputed health with a hardcoded None provider, so it could only ever report a disabled-based reason. - Provider-contract traits (Generation/Page/PageIterator/KeyLookup): provider dispatch went the ListObjectsIndexProviderKind enum route; only ListMetadataIndexHealth is live (dyn in source-mode selection) and stays, minus its unused is_healthy default. - Delete the unused list_quorum_from_env/RUSTFS_API_LIST_QUORUM pair, fold should_resume_local_decommission and should_auto_start_rebalance_after_recovered_meta into their surviving primitives, drop get_disk_via_endpoint, and cfg(test) the remaining test-only conveniences. --- crates/ecstore/src/store/init.rs | 66 ++--- crates/ecstore/src/store/list_objects.rs | 303 ++--------------------- crates/ecstore/src/store/mod.rs | 2 - crates/ecstore/src/store/peer.rs | 4 - 4 files changed, 41 insertions(+), 334 deletions(-) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index bc566d855..1c41694eb 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -49,21 +49,6 @@ fn resolve_startup_pool_defaults_with( drive_counts.into_iter().map(ec_drives_no_config).collect() } -fn should_resume_local_decommission(endpoints: &EndpointServerPools, idx: usize) -> Result { - let pool = endpoints.as_ref().get(idx).ok_or_else(|| { - Error::other(format!( - "store init failed to resolve decommission resume pool index {idx} from current endpoints" - )) - })?; - let endpoint = pool.endpoints.as_ref().first().ok_or_else(|| { - Error::other(format!( - "store init failed to resolve decommission resume pool index {idx}: no endpoints available" - )) - })?; - - Ok(endpoint.is_local) -} - const LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES: usize = 6; const LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY: Duration = Duration::from_secs(60 * 3); const LOCAL_DECOMMISSION_RESUME_RETRY_DELAY: Duration = Duration::from_secs(30); @@ -84,10 +69,6 @@ fn pool_meta_has_active_decommission(meta: &PoolMeta) -> bool { }) } -fn should_auto_start_rebalance_after_recovered_meta(pool_meta: &PoolMeta, rebalance_meta_loaded: bool) -> bool { - should_auto_start_rebalance_after_init(pool_meta_has_active_decommission(pool_meta), rebalance_meta_loaded) -} - async fn wait_for_local_decommission_resume_delay(rx: &CancellationToken, delay: Duration) -> bool { tokio::select! { _ = rx.cancelled() => false, @@ -522,9 +503,9 @@ impl ECStore { mod tests { use super::{ LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, load_pool_meta_for_startup, pool_first_endpoint_is_local, - resolve_startup_pool_defaults_with, resolve_store_init_stage_result, save_validated_pool_meta_for_startup, - should_auto_start_rebalance_after_init, should_auto_start_rebalance_after_recovered_meta, - should_resume_local_decommission, should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay, + pool_meta_has_active_decommission, resolve_startup_pool_defaults_with, resolve_store_init_stage_result, + save_validated_pool_meta_for_startup, should_auto_start_rebalance_after_init, should_retry_local_decommission_resume, + wait_for_local_decommission_resume_delay, }; #[cfg(feature = "test-util")] use crate::{ @@ -696,46 +677,33 @@ mod tests { } #[test] - fn test_should_resume_local_decommission_respects_local_flag() { + fn test_pool_first_endpoint_is_local_respects_local_flag() { let mut local_endpoint = Endpoint::try_from("http://127.0.0.1:9000/data").expect("endpoint should parse"); local_endpoint.is_local = true; - let endpoints = EndpointServerPools::from(vec![PoolEndpoints { + let pool = PoolEndpoints { legacy: false, set_count: 1, drives_per_set: 1, endpoints: Endpoints::from(vec![local_endpoint]), cmd_line: "pool-0".to_string(), platform: String::new(), - }]); + }; - assert!(should_resume_local_decommission(&endpoints, 0).expect("local endpoint should resume")); + assert!(pool_first_endpoint_is_local(&pool)); } #[test] - fn test_should_resume_local_decommission_rejects_unresolvable_pool() { - let endpoints = EndpointServerPools::default(); - let err = should_resume_local_decommission(&endpoints, 0).expect_err("missing pool should error"); - assert_eq!( - err.to_string(), - "Io error: store init failed to resolve decommission resume pool index 0 from current endpoints" - ); - } - - #[test] - fn test_should_resume_local_decommission_rejects_missing_endpoint() { - let endpoints = EndpointServerPools::from(vec![PoolEndpoints { + fn test_pool_first_endpoint_is_local_rejects_missing_endpoint() { + let pool = PoolEndpoints { legacy: false, set_count: 1, drives_per_set: 1, endpoints: Endpoints::from(Vec::::new()), cmd_line: "pool-0".to_string(), platform: String::new(), - }]); - let err = should_resume_local_decommission(&endpoints, 0).expect_err("missing endpoint should error"); - assert_eq!( - err.to_string(), - "Io error: store init failed to resolve decommission resume pool index 0: no endpoints available" - ); + }; + + assert!(!pool_first_endpoint_is_local(&pool)); } #[test] @@ -782,7 +750,10 @@ mod tests { })); let rebalance_meta = Some(RebalanceMeta::default()); - assert!(!should_auto_start_rebalance_after_recovered_meta(&pool_meta, rebalance_meta.is_some())); + assert!(!should_auto_start_rebalance_after_init( + pool_meta_has_active_decommission(&pool_meta), + rebalance_meta.is_some() + )); } #[test] @@ -790,7 +761,10 @@ mod tests { let pool_meta = init_test_pool_meta(None); let rebalance_meta = Some(RebalanceMeta::default()); - assert!(should_auto_start_rebalance_after_recovered_meta(&pool_meta, rebalance_meta.is_some())); + assert!(should_auto_start_rebalance_after_init( + pool_meta_has_active_decommission(&pool_meta), + rebalance_meta.is_some() + )); } #[test] diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 17ce6dae8..758b6279c 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -301,7 +301,6 @@ const LIST_CURSOR_SOURCE_INDEX_KEY_ONLY: &str = "index_key_only"; const LIST_CURSOR_SOURCE_INDEX_VERIFIED_PAGE: &str = "index_verified_page"; const LIST_CURSOR_SOURCE_INDEX_METADATA_FAST: &str = "index_metadata_fast"; const LIST_CURSOR_GENERATION_LIVE: &str = "live"; -const ENV_API_LIST_QUORUM: &str = "RUSTFS_API_LIST_QUORUM"; const DEFAULT_API_LIST_QUORUM: &str = "strict"; const ENV_API_LIST_OBJECTS_QUORUM: &str = "RUSTFS_LIST_OBJECTS_QUORUM"; const DEFAULT_API_LIST_OBJECTS_QUORUM: &str = "optimal"; @@ -406,7 +405,6 @@ enum ListIndexFallbackReason { Unhealthy, Lagging, Degraded, - Corrupt, MissingGeneration, GenerationMismatch, UnsupportedRequest, @@ -421,7 +419,6 @@ impl ListIndexFallbackReason { Self::Unhealthy => "unhealthy", Self::Lagging => "lagging", Self::Degraded => "degraded", - Self::Corrupt => "corrupt", Self::MissingGeneration => "missing_generation", Self::GenerationMismatch => "generation_mismatch", Self::UnsupportedRequest => "unsupported_request", @@ -492,13 +489,6 @@ impl ListObjectsIndexProviderState { } } - fn from_kind(kind: ListObjectsIndexProviderKind) -> Self { - match kind { - ListObjectsIndexProviderKind::WalkerKeyOnly => Self::walker_key_only(), - ListObjectsIndexProviderKind::PersistentKeyOnly => Self::persistent_key_only(None, None), - } - } - fn health_snapshot(&self) -> ListIndexHealthSnapshot { self.lifecycle.health_snapshot() } @@ -667,6 +657,9 @@ async fn observe_list_objects_mutations_with_store(store: Option<&ECStore>, buck Some(next) } +// Test-only observer over the live mutation snapshot; production reads go +// through `current_list_objects_mutation_snapshot` directly. +#[cfg(test)] async fn current_list_objects_mutation_sequence(bucket: &str) -> u64 { current_list_objects_mutation_snapshot(None, bucket).await.high_water_mark } @@ -1539,17 +1532,6 @@ fn generated_persistent_key_only_generation(configured: Option<&str>) -> String format!("{LIST_OBJECTS_INDEX_PROVIDER_PERSISTENT_KEY_ONLY_DEFAULT_GENERATION}-{millis}") } -async fn write_persistent_key_only_index( - store: Option<&ECStore>, - path: &Path, - bucket: &str, - generation: &str, - checkpoint_high_water_mark: u64, - keys: &[String], -) -> Result { - write_persistent_key_only_index_with_metadata(store, path, bucket, generation, checkpoint_high_water_mark, keys, &[]).await -} - async fn write_persistent_key_only_index_with_metadata( store: Option<&ECStore>, path: &Path, @@ -1656,38 +1638,8 @@ enum ListIndexSourceDecision { FallbackToWalker(ListIndexFallbackReason), } -trait ListMetadataIndexGeneration { - fn generation_id(&self) -> &str; -} - trait ListMetadataIndexHealth { fn fallback_reason(&self) -> Option; - - fn is_healthy(&self) -> bool { - self.fallback_reason().is_none() - } -} - -trait ListMetadataIndexPage { - fn source_mode(&self) -> ListSourceMode; - fn generation(&self) -> &dyn ListMetadataIndexGeneration; - fn keys(&self) -> &[String]; - - fn metadata_authority(&self) -> ListMetadataAuthority { - self.source_mode().metadata_authority() - } -} - -trait ListMetadataIndexPageIterator { - type Page: ListMetadataIndexPage; - - fn next_page(&mut self) -> Option; -} - -trait ListMetadataIndexKeyLookup { - type Page: ListMetadataIndexPage; - - fn lookup_page(&self, opts: &ListPathOptions) -> ListIndexSourceDecision; } fn select_list_index_provider_source_mode( @@ -1748,7 +1700,6 @@ enum ListIndexLifecycleState { Healthy, Lagging, Degraded, - Corrupt, } impl ListIndexLifecycleState { @@ -1759,7 +1710,6 @@ impl ListIndexLifecycleState { Self::Healthy => None, Self::Lagging => Some(ListIndexFallbackReason::Lagging), Self::Degraded => Some(ListIndexFallbackReason::Degraded), - Self::Corrupt => Some(ListIndexFallbackReason::Corrupt), } } } @@ -1826,16 +1776,6 @@ impl ListIndexLifecycle { true } - fn recover_after_restart(mut self) -> Self { - self.staging_generation = None; - self.state = if self.active_generation.is_some() { - ListIndexLifecycleState::Healthy - } else { - ListIndexLifecycleState::Disabled - }; - self - } - fn observe_mutation_high_water_mark(&mut self, mutation_high_water_mark: u64) { self.mutation_high_water_mark = self.mutation_high_water_mark.max(mutation_high_water_mark); if matches!(self.state, ListIndexLifecycleState::Healthy | ListIndexLifecycleState::Lagging) @@ -1849,10 +1789,6 @@ impl ListIndexLifecycle { self.state = ListIndexLifecycleState::Degraded; } - fn mark_corrupt(&mut self) { - self.state = ListIndexLifecycleState::Corrupt; - } - fn mutation_lag(&self) -> u64 { self.mutation_high_water_mark.saturating_sub(self.checkpoint_high_water_mark) } @@ -2003,11 +1939,6 @@ fn normalize_list_quorum(value: &str) -> &'static str { } } -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_objects_quorum_from_env() -> String { let value = rustfs_utils::get_env_str(ENV_API_LIST_OBJECTS_QUORUM, DEFAULT_API_LIST_OBJECTS_QUORUM); normalize_list_quorum(&value).to_owned() @@ -2086,38 +2017,10 @@ fn list_objects_index_provider_state_from_env() -> Option) -> ListIndexHealthSnapshot { - provider - .map(ListObjectsIndexProviderState::from_kind) - .map(|state| state.health_snapshot()) - .unwrap_or_else(|| ListIndexLifecycle::disabled(0).health_snapshot()) -} - fn record_list_objects_index_fallback(mode: ListSourceMode, reason: ListIndexFallbackReason) { rustfs_io_metrics::record_list_objects_index_fallback(mode.cursor_value(), reason.metric_label()); } -fn record_list_objects_index_opt_in_fallback(opts: &ListPathOptions, mode: ListSourceMode) { - if !rustfs_io_metrics::get_stage_metrics_enabled() { - return; - } - - let health = list_objects_key_only_provider_health(None); - let reason = match select_list_index_provider_source_mode(opts, mode, &health) { - ListIndexSourceDecision::FallbackToWalker(reason) => reason, - ListIndexSourceDecision::UseIndex(_) => ListIndexFallbackReason::UnsupportedRequest, - }; - - record_list_objects_index_fallback(mode, reason); - debug!( - bucket = %opts.bucket, - prefix = %opts.prefix, - source = mode.cursor_value(), - reason = reason.metric_label(), - "list_objects opt-in index path fell back to walker" - ); -} - fn append_list_cache_id_to_marker(marker: String, cache_id: Option<&str>) -> String { let Some(id) = cache_id else { return marker; @@ -2284,6 +2187,8 @@ struct VerifiedIndexCandidateResult { stats: VerifiedIndexCandidateStats, } +// Test-only convenience over the live `_with_optional_stats` implementation. +#[cfg(test)] async fn list_objects_from_verified_index_candidates( prefix: &str, marker: Option<&str>, @@ -2309,6 +2214,8 @@ where .info) } +// Test-only convenience over the live `_with_optional_stats` implementation. +#[cfg(test)] async fn list_objects_from_verified_index_candidates_with_stats( prefix: &str, marker: Option<&str>, @@ -6753,11 +6660,10 @@ mod test { ENV_API_LIST_OBJECTS_METADATA_FAST_STALENESS_MS, ENV_API_LIST_OBJECTS_NAMESPACE_JOURNAL_CHAOS_BUCKET, ENV_API_LIST_OBJECTS_NAMESPACE_JOURNAL_CHAOS_ENABLED, ENV_API_LIST_OBJECTS_NAMESPACE_JOURNAL_CHAOS_SEQUENCE, ENV_API_LIST_OBJECTS_NAMESPACE_JOURNAL_CHAOS_STATUS, ENV_API_LIST_OBJECTS_NAMESPACE_JOURNAL_PATH, - ENV_API_LIST_OBJECTS_QUORUM, ENV_API_LIST_QUORUM, FallbackListingEntries, FallbackListingEntry, GatherResultsState, + ENV_API_LIST_OBJECTS_QUORUM, FallbackListingEntries, FallbackListingEntry, GatherResultsState, LIST_CURSOR_GENERATION_LIVE, LIST_OBJECTS_INDEX_PROVIDER_PERSISTENT_KEY_ONLY, LIST_OBJECTS_INDEX_PROVIDER_WALKER_KEY_ONLY, ListIndexFallbackReason, ListIndexLifecycle, ListIndexLifecycleState, - ListIndexSourceDecision, ListMetadataAuthority, ListMetadataIndexGeneration, ListMetadataIndexHealth, - ListMetadataIndexKeyLookup, ListMetadataIndexPage, ListMetadataIndexPageIterator, ListObjectsIndexProviderKind, + ListIndexSourceDecision, ListMetadataAuthority, ListMetadataIndexHealth, ListObjectsIndexProviderKind, ListObjectsIndexProviderState, ListPathOptions, ListPathRawOptions, ListSourceMode, ListingEntryResolution, ListingSupplement, ListingSupplementOptions, MAX_OBJECT_LIST, NamespaceMutationJournalBackend, NamespaceMutationJournalSnapshot, NamespaceMutationJournalStatus, PERSISTENT_KEY_ONLY_INDEX_BUCKET_HEADER, @@ -6770,19 +6676,19 @@ mod test { list_merged_entry_channel, list_metadata_resolution_params, list_objects_from_metadata_snapshot_candidates, list_objects_from_verified_index_candidates, list_objects_from_verified_index_candidates_with_optional_stats, list_objects_from_verified_index_candidates_with_stats, list_objects_index_mode_from_env, - list_objects_index_provider_from_env, list_objects_index_provider_state_from_env, list_objects_key_only_provider_health, + list_objects_index_provider_from_env, list_objects_index_provider_state_from_env, list_objects_metadata_fast_guardrails_from_env, list_objects_paginate, list_objects_quorum_from_env, - list_quorum_from_env, load_namespace_mutation_journal_state, load_persistent_key_only_index, max_keys_plus_one, - merge_entry_channels, namespace_mutation_journal_chaos_bucket_from_env, namespace_mutation_journal_chaos_config_from_env, + load_namespace_mutation_journal_state, load_persistent_key_only_index, max_keys_plus_one, merge_entry_channels, + namespace_mutation_journal_chaos_bucket_from_env, namespace_mutation_journal_chaos_config_from_env, namespace_mutation_journal_chaos_enabled_from_env, namespace_mutation_journal_chaos_sequence_from_env, namespace_mutation_journal_chaos_status_from_env, normalize_list_quorum, observe_list_objects_mutations_with_store, parse_namespace_mutation_journal_state, parse_persistent_key_only_index, parse_persistent_list_metadata_object, parse_version_marker, persist_observed_list_objects_mutation, persistent_key_only_index_has_complete_metadata_snapshot, - persistent_key_only_index_health, persistent_key_only_index_matches_provider, record_list_objects_index_opt_in_fallback, + persistent_key_only_index_health, persistent_key_only_index_matches_provider, reset_list_objects_mutation_sequences_for_test, resolve_agreed_listing_entry, resolve_listing_entries, scanner_namespace_mutation_generation, select_list_index_provider_source_mode, select_list_index_source_mode, send_or_cancel, version_marker_for_entries, walk_result_from_set_errors, write_namespace_mutation_journal_state, - write_persistent_key_only_index, + write_persistent_key_only_index_with_metadata, }; use crate::cache_value::metacache_set::{FallbackClaimTracker, TestReaderBehavior, list_path_raw}; use crate::disk::{DiskAPI, DiskOption, endpoint::Endpoint, error::DiskError, new_disk}; @@ -6846,16 +6752,6 @@ mod test { } use uuid::Uuid; - struct TestIndexGeneration { - id: String, - } - - impl ListMetadataIndexGeneration for TestIndexGeneration { - fn generation_id(&self) -> &str { - &self.id - } - } - struct TestIndexHealth { reason: Option, } @@ -6866,50 +6762,6 @@ mod test { } } - struct TestIndexPage { - mode: ListSourceMode, - generation: TestIndexGeneration, - keys: Vec, - } - - impl ListMetadataIndexPage for TestIndexPage { - fn source_mode(&self) -> ListSourceMode { - self.mode - } - - fn generation(&self) -> &dyn ListMetadataIndexGeneration { - &self.generation - } - - fn keys(&self) -> &[String] { - &self.keys - } - } - - struct TestIndexIterator { - pages: std::vec::IntoIter, - } - - impl ListMetadataIndexPageIterator for TestIndexIterator { - type Page = TestIndexPage; - - fn next_page(&mut self) -> Option { - self.pages.next() - } - } - - struct TestIndexLookup { - decision: ListIndexSourceDecision, - } - - impl ListMetadataIndexKeyLookup for TestIndexLookup { - type Page = TestIndexPage; - - fn lookup_page(&self, _opts: &ListPathOptions) -> ListIndexSourceDecision { - self.decision - } - } - fn test_meta_entry(name: &str) -> MetaCacheEntry { MetaCacheEntry { name: name.to_owned(), @@ -7894,30 +7746,6 @@ mod test { 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] #[serial_test::serial] fn list_objects_quorum_from_env_defaults_to_optimal() { @@ -8154,9 +7982,7 @@ mod test { fn list_objects_index_provider_from_env_is_default_off() { temp_env::with_var_unset(ENV_API_LIST_OBJECTS_INDEX_PROVIDER, || { assert_eq!(list_objects_index_provider_from_env(), None); - let health = list_objects_key_only_provider_health(list_objects_index_provider_from_env()); - assert_eq!(health.state, ListIndexLifecycleState::Disabled); - assert_eq!(health.fallback_reason, Some(ListIndexFallbackReason::Disabled)); + assert_eq!(list_objects_index_provider_state_from_env(), None); }); } @@ -8168,7 +7994,9 @@ mod test { Some(LIST_OBJECTS_INDEX_PROVIDER_WALKER_KEY_ONLY), || { assert_eq!(list_objects_index_provider_from_env(), Some(ListObjectsIndexProviderKind::WalkerKeyOnly)); - let health = list_objects_key_only_provider_health(list_objects_index_provider_from_env()); + let health = list_objects_index_provider_state_from_env() + .expect("walker provider should be configured") + .health_snapshot(); assert_eq!(health.state, ListIndexLifecycleState::Healthy); assert_eq!(health.active_generation.as_deref(), Some(LIST_CURSOR_GENERATION_LIVE)); assert_eq!(health.fallback_reason, None); @@ -8459,7 +8287,7 @@ mod test { let journal_root = tempdir.path().join("namespace-mutation-journal"); let keys = vec!["object-a".to_string(), "object-b".to_string()]; - write_persistent_key_only_index(None, &index_path, "bucket", "generation-42", 5, &keys) + write_persistent_key_only_index_with_metadata(None, &index_path, "bucket", "generation-42", 5, &keys, &[]) .await .expect("index should be written"); write_namespace_mutation_journal_state(NamespaceMutationJournalBackend::LocalPath(journal_root), "bucket", 9, false) @@ -8644,7 +8472,7 @@ mod test { #[test] fn list_objects_index_provider_state_uses_lifecycle_active_generation() { - let provider = ListObjectsIndexProviderState::from_kind(ListObjectsIndexProviderKind::WalkerKeyOnly); + let provider = ListObjectsIndexProviderState::walker_key_only(); let health = provider.health_snapshot(); assert_eq!(provider.kind, ListObjectsIndexProviderKind::WalkerKeyOnly); @@ -8657,21 +8485,6 @@ mod test { ); } - #[test] - fn list_objects_index_opt_in_fallback_accepts_incompatible_cursor_without_panicking() { - record_list_objects_index_opt_in_fallback( - &ListPathOptions { - bucket: "bucket".to_string(), - prefix: "photos/".to_string(), - marker: Some( - "photos/2026/[rustfs_cache:v2,id:list-cache-id,src:index_verified_page,gen:generation-1]".to_string(), - ), - ..Default::default() - }, - ListSourceMode::IndexKeyOnly, - ); - } - #[tokio::test] async fn list_objects_verified_index_candidates_use_live_metadata_only() { let candidates = vec![ @@ -9392,38 +9205,6 @@ mod test { ); } - #[test] - fn list_index_traits_keep_page_iteration_and_metadata_authority_explicit() { - let page = TestIndexPage { - mode: ListSourceMode::IndexMetadataFast, - generation: TestIndexGeneration { - id: "snapshot-7".to_string(), - }, - keys: vec!["photos/2026/image.jpg".to_string()], - }; - let mut iterator = TestIndexIterator { - pages: vec![page].into_iter(), - }; - let page = iterator.next_page().expect("test iterator should produce one page"); - - assert_eq!(page.keys(), &["photos/2026/image.jpg".to_string()]); - assert_eq!(page.generation().generation_id(), "snapshot-7"); - assert_eq!(page.metadata_authority(), ListMetadataAuthority::IndexSnapshotEventuallyConsistent); - assert!(!page.source_mode().can_satisfy_strong_listing()); - } - - #[test] - fn list_index_key_lookup_trait_returns_explicit_fallback_decision() { - let lookup = TestIndexLookup { - decision: ListIndexSourceDecision::FallbackToWalker(ListIndexFallbackReason::UnsupportedRequest), - }; - - assert_eq!( - lookup.lookup_page(&ListPathOptions::default()), - ListIndexSourceDecision::FallbackToWalker(ListIndexFallbackReason::UnsupportedRequest) - ); - } - #[test] fn list_index_lifecycle_rebuild_requires_publish_before_health() { let mut lifecycle = ListIndexLifecycle::disabled(10); @@ -9446,40 +9227,6 @@ mod test { assert_eq!(snapshot.fallback_reason, None); } - #[test] - fn list_index_lifecycle_restart_never_trusts_unpublished_staging_generation() { - let mut lifecycle = ListIndexLifecycle::disabled(10); - lifecycle.begin_rebuild("generation-1", 100); - lifecycle.checkpoint_rebuild(90); - - let recovered = lifecycle.recover_after_restart(); - let snapshot = recovered.health_snapshot(); - - assert_eq!(snapshot.state, ListIndexLifecycleState::Disabled); - assert_eq!(snapshot.active_generation, None); - assert_eq!(snapshot.staging_generation, None); - assert_eq!(snapshot.fallback_reason, Some(ListIndexFallbackReason::Disabled)); - } - - #[test] - fn list_index_lifecycle_restart_keeps_only_last_published_generation() { - let mut lifecycle = ListIndexLifecycle::disabled(10); - lifecycle.begin_rebuild("generation-1", 100); - lifecycle.checkpoint_rebuild(100); - assert!(lifecycle.publish_rebuild()); - - lifecycle.begin_rebuild("generation-2", 125); - lifecycle.checkpoint_rebuild(120); - - let recovered = lifecycle.recover_after_restart(); - let snapshot = recovered.health_snapshot(); - - assert_eq!(snapshot.state, ListIndexLifecycleState::Healthy); - assert_eq!(snapshot.active_generation.as_deref(), Some("generation-1")); - assert_eq!(snapshot.staging_generation, None); - assert_eq!(snapshot.fallback_reason, None); - } - #[test] fn list_index_lifecycle_reports_lagging_and_fallback_reason() { let mut lifecycle = ListIndexLifecycle::disabled(10); @@ -9493,11 +9240,10 @@ mod test { assert_eq!(snapshot.state, ListIndexLifecycleState::Lagging); assert_eq!(snapshot.mutation_lag, 15); assert_eq!(snapshot.fallback_reason, Some(ListIndexFallbackReason::Lagging)); - assert!(!snapshot.is_healthy()); } #[test] - fn list_index_lifecycle_reports_degraded_and_corrupt_as_unhealthy() { + fn list_index_lifecycle_reports_degraded_as_fallback() { let mut lifecycle = ListIndexLifecycle::disabled(10); lifecycle.begin_rebuild("generation-1", 100); lifecycle.checkpoint_rebuild(100); @@ -9507,13 +9253,6 @@ mod test { let degraded = lifecycle.health_snapshot(); assert_eq!(degraded.state, ListIndexLifecycleState::Degraded); assert_eq!(degraded.fallback_reason, Some(ListIndexFallbackReason::Degraded)); - assert!(!degraded.is_healthy()); - - lifecycle.mark_corrupt(); - let corrupt = lifecycle.health_snapshot(); - assert_eq!(corrupt.state, ListIndexLifecycleState::Corrupt); - assert_eq!(corrupt.fallback_reason, Some(ListIndexFallbackReason::Corrupt)); - assert!(!corrupt.is_healthy()); } #[test] diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index f706aef15..bad794f01 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -12,8 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -// #730: store startup/listing helpers are staged while global state is moved to owners. -#![allow(dead_code)] #![allow(clippy::map_entry)] use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc; diff --git a/crates/ecstore/src/store/peer.rs b/crates/ecstore/src/store/peer.rs index 29919b6d0..e827c7fa4 100644 --- a/crates/ecstore/src/store/peer.rs +++ b/crates/ecstore/src/store/peer.rs @@ -61,10 +61,6 @@ pub async fn find_local_disk_by_ref(disk_ref: &str) -> Option { None } -pub async fn get_disk_via_endpoint(endpoint: &Endpoint) -> Option { - runtime_sources::local_disk_for_endpoint(endpoint).await -} - pub async fn all_local_disk_path() -> Vec { runtime_sources::local_disk_paths().await }