diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 860ebb2cd..f568050d8 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -105,6 +105,12 @@ struct DirtyUsageSnapshot { covers_all_pending: bool, } +#[derive(Clone, Debug, Default)] +pub(crate) struct ScannerBucketScanScope { + selected_buckets: Option>>, + baseline_scan_plan_digest: Option, +} + pub(crate) fn is_scanner_metadata_corrupt_error(err: &StorageError) -> bool { matches!(err, StorageError::Io(io) if io.to_string().starts_with(SCANNER_METADATA_CORRUPT_ERROR)) } @@ -146,6 +152,7 @@ fn object_lock_config_enabled(config: &ObjectLockConfiguration) -> bool { pub struct ScannerBucketScanPlan { buckets: Vec, all_buckets: Arc>, + scope: ScannerBucketScanScope, digest: DataUsageScanPlanDigest, leader_epoch: u64, tier_registry_generation: u64, @@ -732,6 +739,8 @@ mod dirty_usage; mod guards; mod io_cache; mod io_cycle; +#[cfg(test)] +use io_cache::{ScannerSetCacheGeneration, prepare_scoped_set_scan}; pub(crate) use io_cycle::nsscanner_with_storage_status; mod io_disk; #[cfg(test)] diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index c99ffacd1..a8686a1e4 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -14,6 +14,93 @@ /// ScannerIOCache implementation for SetDisks: bucket ordering, worker fan-out, merge, and publish. use super::*; +#[derive(Clone, Copy)] +pub(super) struct ScannerSetCacheGeneration { + pub(super) want_cycle: u64, + pub(super) leader_epoch: u64, + pub(super) tier_registry_generation: u64, + pub(super) source: DataUsageCacheSource, + pub(super) scan_plan_digest: DataUsageScanPlanDigest, +} + +pub(super) struct PreparedScopedSetScan { + pub(super) buckets: Vec, + pub(super) cache: DataUsageCache, +} + +pub(super) fn prepare_scoped_set_scan( + old_cache: &DataUsageCache, + set_buckets: &[BucketInfo], + all_buckets: &[BucketInfo], + scope: &ScannerBucketScanScope, + generation: ScannerSetCacheGeneration, +) -> Option { + let (Some(selected_buckets), Some(baseline_scan_plan_digest)) = (&scope.selected_buckets, scope.baseline_scan_plan_digest) + else { + return None; + }; + if selected_buckets.is_empty() + || !old_cache.info.snapshot_complete + || old_cache.info.last_update.is_none() + || old_cache.info.name != DATA_USAGE_ROOT + || old_cache.info.next_cycle > generation.want_cycle + || old_cache.info.leader_epoch != generation.leader_epoch + || old_cache.info.tier_registry_generation != Some(generation.tier_registry_generation) + || old_cache.info.source != Some(generation.source) + || old_cache.info.scan_plan_digest != Some(baseline_scan_plan_digest) + || old_cache.info.cache_key_format != DATA_USAGE_CACHE_KEY_FORMAT + || old_cache.checked_flatten_complete_scope(DATA_USAGE_ROOT).is_none() + { + return None; + } + + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + next_cycle: generation.want_cycle, + leader_epoch: generation.leader_epoch, + tier_registry_generation: Some(generation.tier_registry_generation), + source: Some(generation.source), + snapshot_complete: false, + scan_plan_digest: Some(generation.scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + lkg_snapshot_complete: true, + lkg_next_cycle: Some(old_cache.info.next_cycle), + lkg_last_update: old_cache.info.last_update, + lkg_leader_epoch: Some(old_cache.info.leader_epoch), + lkg_scan_plan_digest: old_cache.info.scan_plan_digest, + ..Default::default() + }, + cache: HashMap::new(), + }; + cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); + let root_hash = crate::hash_path(DATA_USAGE_ROOT); + let mut current_bucket_names = HashSet::with_capacity(all_buckets.len()); + for bucket in all_buckets { + if !current_bucket_names.insert(bucket.name.as_str()) { + return None; + } + if selected_buckets.contains(&bucket.name) { + cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default()); + continue; + } + + let bucket_hash = crate::hash_path(&bucket.name); + old_cache.find(&bucket.name)?; + cache.copy_with_children(old_cache, &bucket_hash, &Some(root_hash.clone())); + cache.find(&bucket.name)?; + } + + Some(PreparedScopedSetScan { + buckets: set_buckets + .iter() + .filter(|bucket| selected_buckets.contains(&bucket.name)) + .cloned() + .collect(), + cache, + }) +} + #[async_trait::async_trait] impl ScannerIOCache for SetDisks { #[tracing::instrument(skip(self, budget, scan_plan, updates))] @@ -27,8 +114,9 @@ impl ScannerIOCache for SetDisks { scan_mode: HealScanMode, ) -> Result<()> { let ScannerBucketScanPlan { - buckets, + mut buckets, all_buckets, + scope, digest: scan_plan_digest, leader_epoch, tier_registry_generation, @@ -63,26 +151,57 @@ impl ScannerIOCache for SetDisks { "Scanner old data usage cache load failed; rebuilding from bucket caches" ); } + let scoped_scan = prepare_scoped_set_scan( + &old_cache, + &buckets, + &all_buckets, + &scope, + ScannerSetCacheGeneration { + want_cycle, + leader_epoch, + tier_registry_generation, + source, + scan_plan_digest, + }, + ); + let mut scoped_cache = scoped_scan.map(|prepared| { + buckets = prepared.buckets; + prepared.cache + }); if buckets.is_empty() { let now = SystemTime::now(); - let mut cache = DataUsageCache { - info: DataUsageCacheInfo { - name: DATA_USAGE_ROOT.to_string(), - next_cycle: want_cycle, - last_update: Some(now), - leader_epoch, - tier_registry_generation: Some(tier_registry_generation), - source: Some(source), - snapshot_complete: true, - scan_plan_digest: Some(scan_plan_digest), - cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, - ..Default::default() - }, - cache: HashMap::new(), + let mut cache = match scoped_cache.take() { + Some(cache) => cache, + None => { + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + next_cycle: want_cycle, + leader_epoch, + tier_registry_generation: Some(tier_registry_generation), + source: Some(source), + scan_plan_digest: Some(scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + cache: HashMap::new(), + }; + cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); + for bucket in all_buckets.iter() { + cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default()); + } + cache + } }; - cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); - for bucket in all_buckets.iter() { - cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default()); + cache.info.last_update = Some(now); + cache.info.snapshot_complete = true; + cache.info.lkg_snapshot_complete = false; + cache.info.lkg_next_cycle = None; + cache.info.lkg_last_update = None; + cache.info.lkg_leader_epoch = None; + cache.info.lkg_scan_plan_digest = None; + if cache.find(DATA_USAGE_ROOT).is_none() { + cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); } reset_disk_bucket_scan_gauges(&pool_label, &set_label); return persist_and_publish_cache_snapshot( @@ -269,92 +388,102 @@ impl ScannerIOCache for SetDisks { record_disk_bucket_scans_active(0, &pool_label, &set_label); let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone()); - // Fence a stale set aggregate before copying entries into per-bucket work caches. - if old_cache.info.next_cycle <= want_cycle - && old_cache.info.leader_epoch <= leader_epoch - && old_cache.info.tier_registry_generation != Some(tier_registry_generation) - { - old_cache.info.scan_plan_digest = None; - } - let old_lkg = old_cache.info.snapshot_complete.then_some({ - ( - old_cache.info.next_cycle, - old_cache.info.last_update, - old_cache.info.leader_epoch, - old_cache.info.scan_plan_digest, - ) - }); - let prepare_outcome = match old_cache.prepare_for_scan( - DATA_USAGE_ROOT, - want_cycle, - leader_epoch, - source, - scan_plan_digest, - require_cache_source, - ) { - DataUsageCachePrepareOutcome::RejectedNewerCycle => { - cache_cycle_floor.fetch_max(old_cache.info.next_cycle, Ordering::AcqRel); - warn!( - target: "rustfs::scanner::io", - event = EVENT_SCANNER_CACHE_PERSIST_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_IO, - pool = self.pool_index, - set = self.set_index, - cache_name = DATA_USAGE_CACHE_NAME, - requested_cycle = want_cycle, - cached_cycle = old_cache.info.next_cycle, - state = "stale_cycle_rejected", - "Scanner rejected a set cache cycle regression" - ); - return Ok(()); + let mut cache = if let Some(cache) = scoped_cache.take() { + cache + } else { + // Fence a stale set aggregate before copying entries into per-bucket work caches. + if old_cache.info.next_cycle <= want_cycle + && old_cache.info.leader_epoch <= leader_epoch + && old_cache.info.tier_registry_generation != Some(tier_registry_generation) + { + old_cache.info.scan_plan_digest = None; } - DataUsageCachePrepareOutcome::RejectedNewerLeader => { - warn!( - target: "rustfs::scanner::io", - event = EVENT_SCANNER_CACHE_PERSIST_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_IO, - pool = self.pool_index, - set = self.set_index, - cache_name = DATA_USAGE_CACHE_NAME, - requested_epoch = leader_epoch, - cached_epoch = old_cache.info.leader_epoch, - state = "stale_leader_rejected", - "Scanner rejected work from an older leader epoch" - ); - return Ok(()); - } - outcome => outcome, - }; - if matches!(prepare_outcome, DataUsageCachePrepareOutcome::Reused) - && let Some((cycle, last_update, epoch, digest)) = old_lkg - { - old_cache.info.lkg_snapshot_complete = true; - old_cache.info.lkg_next_cycle = Some(cycle); - old_cache.info.lkg_last_update = last_update; - old_cache.info.lkg_leader_epoch = Some(epoch); - old_cache.info.lkg_scan_plan_digest = digest; - } - - let mut cache = DataUsageCache { - info: DataUsageCacheInfo { - name: DATA_USAGE_ROOT.to_string(), - next_cycle: want_cycle, + let old_lkg = old_cache.info.snapshot_complete.then_some({ + ( + old_cache.info.next_cycle, + old_cache.info.last_update, + old_cache.info.leader_epoch, + old_cache.info.scan_plan_digest, + ) + }); + let prepare_outcome = match old_cache.prepare_for_scan( + DATA_USAGE_ROOT, + want_cycle, leader_epoch, - tier_registry_generation: Some(tier_registry_generation), - source: Some(source), - snapshot_complete: false, - scan_plan_digest: Some(scan_plan_digest), - cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, - ..Default::default() - }, - cache: HashMap::new(), + source, + scan_plan_digest, + require_cache_source, + ) { + DataUsageCachePrepareOutcome::RejectedNewerCycle => { + cache_cycle_floor.fetch_max(old_cache.info.next_cycle, Ordering::AcqRel); + warn!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_CACHE_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + pool = self.pool_index, + set = self.set_index, + cache_name = DATA_USAGE_CACHE_NAME, + requested_cycle = want_cycle, + cached_cycle = old_cache.info.next_cycle, + state = "stale_cycle_rejected", + "Scanner rejected a set cache cycle regression" + ); + return Ok(()); + } + DataUsageCachePrepareOutcome::RejectedNewerLeader => { + warn!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_CACHE_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + pool = self.pool_index, + set = self.set_index, + cache_name = DATA_USAGE_CACHE_NAME, + requested_epoch = leader_epoch, + cached_epoch = old_cache.info.leader_epoch, + state = "stale_leader_rejected", + "Scanner rejected work from an older leader epoch" + ); + return Ok(()); + } + outcome => outcome, + }; + if matches!(prepare_outcome, DataUsageCachePrepareOutcome::Reused) + && let Some((cycle, last_update, epoch, digest)) = old_lkg + { + old_cache.info.lkg_snapshot_complete = true; + old_cache.info.lkg_next_cycle = Some(cycle); + old_cache.info.lkg_last_update = last_update; + old_cache.info.lkg_leader_epoch = Some(epoch); + old_cache.info.lkg_scan_plan_digest = digest; + } + + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + next_cycle: want_cycle, + leader_epoch, + tier_registry_generation: Some(tier_registry_generation), + source: Some(source), + snapshot_complete: false, + scan_plan_digest: Some(scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + lkg_snapshot_complete: old_cache.info.lkg_snapshot_complete, + lkg_next_cycle: old_cache.info.lkg_next_cycle, + lkg_last_update: old_cache.info.lkg_last_update, + lkg_leader_epoch: old_cache.info.lkg_leader_epoch, + lkg_scan_plan_digest: old_cache.info.lkg_scan_plan_digest, + ..Default::default() + }, + cache: HashMap::new(), + }; + cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); + for bucket in all_buckets.iter() { + cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default()); + } + cache }; - cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); - for bucket in all_buckets.iter() { - cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default()); - } let (bucket_tx, bucket_rx) = mpsc::channel::(buckets.len()); @@ -1257,11 +1386,6 @@ impl ScannerIOCache for SetDisks { incomplete_scope.info.snapshot_complete = false; incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest); incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT; - incomplete_scope.info.lkg_snapshot_complete = old_cache.info.lkg_snapshot_complete; - incomplete_scope.info.lkg_next_cycle = old_cache.info.lkg_next_cycle; - incomplete_scope.info.lkg_last_update = old_cache.info.lkg_last_update; - incomplete_scope.info.lkg_leader_epoch = old_cache.info.lkg_leader_epoch; - incomplete_scope.info.lkg_scan_plan_digest = old_cache.info.lkg_scan_plan_digest; if let Err(e) = updates.send(incomplete_scope).await { error!( target: "rustfs::scanner::io", diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index f51caf381..3b14e9ff2 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -63,6 +63,41 @@ pub(crate) async fn nsscanner_with_storage_status( where S: ScannerStorage, { + let request = ScannerCycleRequest { + ctx, + budget, + updates, + want_cycle, + leader_epoch, + scan_mode, + scan_scope: ScannerBucketScanScope::default(), + }; + nsscanner_with_storage_status_scoped(store, request).await +} + +pub(crate) struct ScannerCycleRequest { + pub(crate) ctx: CancellationToken, + pub(crate) budget: Arc, + pub(crate) updates: mpsc::Sender, + pub(crate) want_cycle: u64, + pub(crate) leader_epoch: u64, + pub(crate) scan_mode: HealScanMode, + pub(crate) scan_scope: ScannerBucketScanScope, +} + +pub(crate) async fn nsscanner_with_storage_status_scoped(store: &S, request: ScannerCycleRequest) -> Result +where + S: ScannerStorage, +{ + let ScannerCycleRequest { + ctx, + budget, + updates, + want_cycle, + leader_epoch, + scan_mode, + scan_scope, + } = request; let child_token = ctx.child_token(); let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch); @@ -280,6 +315,7 @@ where let scan_plan = ScannerBucketScanPlan { buckets: set_buckets, all_buckets: Arc::clone(&all_buckets), + scope: scan_scope.clone(), digest: scan_plan_digest, leader_epoch, tier_registry_generation, diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 47f1e86a0..a59d620fc 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -765,6 +765,155 @@ fn bucket_usage_scan_order_prioritizes_dirty_buckets() { assert_eq!(names, vec!["dirty", "missing", "cached"]); } +fn complete_set_usage_cache(buckets: &[(&str, usize)], scan_plan_digest: DataUsageScanPlanDigest) -> DataUsageCache { + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + next_cycle: 7, + last_update: Some(SystemTime::now()), + leader_epoch: 11, + source: Some(DataUsageCacheSource::new(1, 2)), + snapshot_complete: true, + scan_plan_digest: Some(scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + tier_registry_generation: Some(13), + ..Default::default() + }, + ..Default::default() + }; + cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); + for (bucket, size) in buckets { + cache.replace( + bucket, + DATA_USAGE_ROOT, + DataUsageEntry { + size: *size, + objects: 1, + ..Default::default() + }, + ); + } + cache +} + +#[test] +fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() { + let baseline_digest = DataUsageScanPlanDigest([1; 32]); + let current_digest = DataUsageScanPlanDigest([2; 32]); + let mut old_cache = complete_set_usage_cache(&[("stable", 10), ("dirty", 20), ("deleted", 30)], baseline_digest); + old_cache.replace( + "stable/prefix", + "stable", + DataUsageEntry { + size: 5, + objects: 1, + ..Default::default() + }, + ); + let all_buckets = vec![bucket_info("stable"), bucket_info("dirty")]; + let selected_buckets = Arc::new(HashSet::from(["dirty".to_string(), "deleted".to_string()])); + + let prepared = prepare_scoped_set_scan( + &old_cache, + &all_buckets, + &all_buckets, + &ScannerBucketScanScope { + selected_buckets: Some(selected_buckets), + baseline_scan_plan_digest: Some(baseline_digest), + }, + ScannerSetCacheGeneration { + want_cycle: 8, + leader_epoch: 11, + tier_registry_generation: 13, + source: DataUsageCacheSource::new(1, 2), + scan_plan_digest: current_digest, + }, + ) + .expect("complete matching set cache should support a scoped scan"); + + assert_eq!(prepared.buckets.iter().map(|bucket| bucket.name.as_str()).collect::>(), ["dirty"]); + let stable = prepared + .cache + .checked_flatten("stable") + .expect("unselected bucket subtree should be retained"); + assert_eq!((stable.size, stable.objects), (15, 2)); + assert_eq!(prepared.cache.find("dirty").map(|entry| (entry.size, entry.objects)), Some((0, 0))); + assert!(prepared.cache.find("deleted").is_none()); + assert_eq!(prepared.cache.info.scan_plan_digest, Some(current_digest)); + assert_eq!(prepared.cache.info.next_cycle, 8); + assert!(!prepared.cache.info.snapshot_complete); + assert!(prepared.cache.info.lkg_snapshot_complete); + assert_eq!(prepared.cache.info.lkg_next_cycle, Some(7)); + assert_eq!(prepared.cache.info.lkg_scan_plan_digest, Some(baseline_digest)); +} + +#[test] +fn scoped_set_scan_falls_back_when_an_unselected_bucket_has_no_baseline() { + let baseline_digest = DataUsageScanPlanDigest([3; 32]); + let old_cache = complete_set_usage_cache(&[("stable", 10)], baseline_digest); + let all_buckets = vec![bucket_info("stable"), bucket_info("new")]; + + assert!( + prepare_scoped_set_scan( + &old_cache, + &all_buckets, + &all_buckets, + &ScannerBucketScanScope { + selected_buckets: Some(Arc::new(HashSet::from(["dirty".to_string()]))), + baseline_scan_plan_digest: Some(baseline_digest), + }, + ScannerSetCacheGeneration { + want_cycle: 8, + leader_epoch: 11, + tier_registry_generation: 13, + source: DataUsageCacheSource::new(1, 2), + scan_plan_digest: DataUsageScanPlanDigest([4; 32]), + }, + ) + .is_none() + ); +} + +#[test] +fn scoped_set_scan_requires_an_exact_complete_baseline() { + let baseline_digest = DataUsageScanPlanDigest([5; 32]); + let all_buckets = vec![bucket_info("dirty")]; + let scope = ScannerBucketScanScope { + selected_buckets: Some(Arc::new(HashSet::from(["dirty".to_string()]))), + baseline_scan_plan_digest: Some(baseline_digest), + }; + let generation = ScannerSetCacheGeneration { + want_cycle: 8, + leader_epoch: 11, + tier_registry_generation: 13, + source: DataUsageCacheSource::new(1, 2), + scan_plan_digest: DataUsageScanPlanDigest([6; 32]), + }; + + let mut incomplete = complete_set_usage_cache(&[("dirty", 10)], baseline_digest); + incomplete.info.snapshot_complete = false; + assert!(prepare_scoped_set_scan(&incomplete, &all_buckets, &all_buckets, &scope, generation).is_none()); + + let mut not_durable = complete_set_usage_cache(&[("dirty", 10)], baseline_digest); + not_durable.info.last_update = None; + assert!(prepare_scoped_set_scan(¬_durable, &all_buckets, &all_buckets, &scope, generation).is_none()); + + let mut wrong_digest = complete_set_usage_cache(&[("dirty", 10)], baseline_digest); + wrong_digest.info.scan_plan_digest = Some(DataUsageScanPlanDigest([7; 32])); + assert!(prepare_scoped_set_scan(&wrong_digest, &all_buckets, &all_buckets, &scope, generation).is_none()); + + let empty_scope = ScannerBucketScanScope { + selected_buckets: Some(Arc::new(HashSet::new())), + baseline_scan_plan_digest: Some(baseline_digest), + }; + let complete = complete_set_usage_cache(&[("dirty", 10)], baseline_digest); + assert!(prepare_scoped_set_scan(&complete, &all_buckets, &all_buckets, &empty_scope, generation).is_none()); + + let mut future_cache = complete_set_usage_cache(&[("dirty", 10)], baseline_digest); + future_cache.info.next_cycle = generation.want_cycle.saturating_add(1); + assert!(prepare_scoped_set_scan(&future_cache, &all_buckets, &all_buckets, &scope, generation).is_none()); +} + #[test] fn record_set_scan_failure_preserves_first_error() { let mut first = None; diff --git a/rustfs/src/version.rs b/rustfs/src/version.rs index d94adedaf..7384b47a7 100644 --- a/rustfs/src/version.rs +++ b/rustfs/src/version.rs @@ -14,8 +14,6 @@ use const_str::concat; use shadow_rs::shadow; -use std::path::Path; -use std::process::Command; shadow!(build); @@ -47,10 +45,6 @@ pub const DISPLAY_VERSION: &str = { type VersionParseResult = Result<(u32, u32, u32, Option), Box>; -fn build_version_override() -> Option<&'static str> { - BUILD_VERSION_OVERRIDE.filter(|version| !version.is_empty()) -} - fn version_ref(version: &str) -> String { if version.starts_with("refs/tags/") || version.starts_with('@') { version.to_string() @@ -61,91 +55,7 @@ fn version_ref(version: &str) -> String { #[allow(clippy::const_is_empty)] pub fn get_version() -> String { - if let Some(version) = build_version_override() { - return version_ref(version); - } - - // Get the latest tag - if let Ok(latest_tag) = get_latest_tag() { - // Check if current commit is newer than the latest tag - if is_head_newer_than_tag(&latest_tag) { - // If current commit is newer, increment the version number - if let Ok(new_version) = increment_version(&latest_tag) { - return format!("refs/tags/{new_version}"); - } - } - - // If current commit is the latest tag, or version increment failed, return current tag - return format!("refs/tags/{latest_tag}"); - } - - // If no tag exists, use original logic - if !build::TAG.is_empty() { - format!("refs/tags/{}", build::TAG) - } else if !build::SHORT_COMMIT.is_empty() { - format!("@{}", build::SHORT_COMMIT) - } else { - format!("refs/tags/{}", build::PKG_VERSION) - } -} - -/// Get the latest git tag -fn get_latest_tag() -> Result> { - let output = Command::new("git").args(["describe", "--tags", "--abbrev=0"]).output()?; - - if output.status.success() { - let tag = String::from_utf8(output.stdout)?; - Ok(tag.trim().to_string()) - } else { - Err("Failed to get latest tag".into()) - } -} - -/// Check if current HEAD is newer than specified tag -fn is_head_newer_than_tag(tag: &str) -> bool { - is_head_newer_than_tag_in(Path::new("."), tag) -} - -fn is_head_newer_than_tag_in(repo: &Path, tag: &str) -> bool { - let head = Command::new("git").current_dir(repo).args(["rev-parse", "HEAD"]).output(); - let tag_commit = Command::new("git") - .current_dir(repo) - .args(["rev-list", "-n", "1", tag]) - .output(); - - let (Ok(head), Ok(tag_commit)) = (head, tag_commit) else { - return false; - }; - - if !head.status.success() || !tag_commit.status.success() || head.stdout == tag_commit.stdout { - return false; - } - - let output = Command::new("git") - .current_dir(repo) - .args(["merge-base", "--is-ancestor", tag, "HEAD"]) - .output(); - - match output { - Ok(result) => result.status.success(), - Err(_) => false, - } -} - -/// Increment version number (increase patch version) -fn increment_version(version: &str) -> Result> { - // Parse version number, e.g. "1.0.0-alpha.19" -> (1, 0, 0, Some("alpha.19")) - let (major, minor, patch, pre_release) = parse_version(version)?; - - // If there's a pre-release identifier, increment the pre-release version number - if let Some(pre) = pre_release - && let Some(new_pre) = increment_pre_release(&pre) - { - return Ok(format!("{major}.{minor}.{patch}-{new_pre}")); - } - - // Otherwise increment patch version number - Ok(format!("{major}.{minor}.{}", patch + 1)) + version_ref(DISPLAY_VERSION) } /// Parse version number @@ -166,28 +76,6 @@ pub fn parse_version(version: &str) -> VersionParseResult { Ok((major, minor, patch, pre_release)) } -/// Increment pre-release version number -fn increment_pre_release(pre_release: &str) -> Option { - // Handle pre-release versions like "alpha.19" - let parts: Vec<&str> = pre_release.split('.').collect(); - if parts.len() == 2 - && let Ok(num) = parts[1].parse::() - { - return Some(format!("{}.{}", parts[0], num + 1)); - } - - // Handle pre-release versions like "alpha19" - if let Some(pos) = pre_release.rfind(|c: char| c.is_alphabetic()) { - let prefix = &pre_release[..=pos]; - let suffix = &pre_release[pos + 1..]; - if let Ok(num) = suffix.parse::() { - return Some(format!("{prefix}{}", num + 1)); - } - } - - None -} - /// Clean version string - removes common prefixes pub fn clean_version(version: &str) -> String { version @@ -284,34 +172,6 @@ mod tests { use super::*; use tracing::debug; - fn run_git(repo: &Path, args: &[&str]) { - let status = Command::new("git").current_dir(repo).args(args).status().unwrap(); - assert!(status.success(), "git command failed: git {}", args.join(" ")); - } - - #[test] - fn test_is_head_newer_than_tag_requires_strict_descendant() { - let repo = tempfile::tempdir().unwrap(); - run_git(repo.path(), &["init", "--quiet"]); - run_git(repo.path(), &["config", "user.name", "RustFS Tests"]); - run_git(repo.path(), &["config", "user.email", "rustfs@example.com"]); - run_git(repo.path(), &["commit", "--allow-empty", "--quiet", "-m", "tagged commit"]); - run_git(repo.path(), &["tag", "--annotate", "1.2.3", "--message", "1.2.3"]); - - assert!(!is_head_newer_than_tag_in(repo.path(), "1.2.3")); - - run_git(repo.path(), &["commit", "--allow-empty", "--quiet", "-m", "newer commit"]); - - assert!(is_head_newer_than_tag_in(repo.path(), "1.2.3")); - } - - #[test] - fn build_version_override_is_used_for_current_version_when_set() { - if let Some(version) = build_version_override() { - assert_eq!(get_version(), version_ref(version)); - } - } - #[test] fn version_ref_keeps_existing_ref_prefixes() { assert_eq!(version_ref("1.2.3"), "refs/tags/1.2.3"); @@ -319,6 +179,11 @@ mod tests { assert_eq!(version_ref("@abc123"), "@abc123"); } + #[test] + fn get_version_uses_build_metadata() { + assert_eq!(get_version(), version_ref(DISPLAY_VERSION)); + } + #[test] fn test_parse_version() { // Test standard version parsing @@ -336,27 +201,6 @@ mod tests { assert_eq!(pre_release, Some("alpha.19".to_string())); } - #[test] - fn test_increment_pre_release() { - // Test alpha.19 -> alpha.20 - assert_eq!(increment_pre_release("alpha.19"), Some("alpha.20".to_string())); - - // Test beta.5 -> beta.6 - assert_eq!(increment_pre_release("beta.5"), Some("beta.6".to_string())); - - // Test unparsable case - assert_eq!(increment_pre_release("unknown"), None); - } - - #[test] - fn test_increment_version() { - // Test pre-release version increment - assert_eq!(increment_version("1.0.0-alpha.19").unwrap(), "1.0.0-alpha.20"); - - // Test standard version increment - assert_eq!(increment_version("1.0.0").unwrap(), "1.0.1"); - } - #[test] fn test_version_format() { // Test if version format starts with refs/tags/