diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 3c32e075b..5b67ff411 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -29,15 +29,15 @@ use storage_api::owner::{ ECSTORE_STORAGECLASS_STANDARD, ECSTORE_TRANSITION_COMPLETE, EcstoreBucketTargetSys, EcstoreBucketVersioningSys, EcstoreDisk, EcstoreDiskAPI, EcstoreDiskBytes, EcstoreDiskError, EcstoreDiskInfo, EcstoreDiskInfoOptions, EcstoreDiskLocation, EcstoreDiskResult, EcstoreErrorType, EcstoreEvaluator, EcstoreEvent, EcstoreLcEventSrc, EcstoreLifecycle, - EcstoreListPathRawOptions, EcstoreObjectOpts, EcstoreReplicationConfig, EcstoreReplicationConfigurationExt, - EcstoreReplicationHealQueueResult, EcstoreReplicationQueueAdmission, EcstoreReplicationScannerBridge, EcstoreResultType, - EcstoreScanGuard, EcstoreSetDisks, EcstoreStorageError, EcstoreStore, EcstoreTierConfig, EcstoreVersioningApi, HTTPRangeSpec, - ObjectIO, ObjectOperations, ObjectToDelete, ecstore_apply_expiry_rule, ecstore_apply_transition_rule, + EcstoreListPathRawOptions, EcstoreObjectOpts, EcstoreReplicationConfigurationExt, EcstoreReplicationScannerBridge, + EcstoreResultType, EcstoreScanGuard, EcstoreSetDisks, EcstoreStorageError, EcstoreStore, EcstoreTierConfig, + EcstoreVersioningApi, HTTPRangeSpec, ObjectIO, ObjectOperations, ObjectToDelete, ScannerReplicationHealObject, + ScannerReplicationHealResult, ScannerReplicationQueueAdmission, ecstore_apply_expiry_rule, ecstore_apply_transition_rule, ecstore_expiry_state_handle, ecstore_get_global_tier_config_mgr, ecstore_get_lifecycle_config, ecstore_get_object_lock_config, ecstore_get_replication_config, ecstore_is_erasure, ecstore_is_erasure_sd, ecstore_is_reserved_or_invalid_bucket, ecstore_list_path_raw, ecstore_path2_bucket_object, ecstore_path2_bucket_object_with_base_path, ecstore_read_config, ecstore_replace_bucket_usage_memory_from_info, - ecstore_resolve_object_store_handle, ecstore_save_config, + ecstore_resolve_object_store_handle, ecstore_save_config, scanner_replication_config_for_lifecycle_eval, }; #[cfg(test)] use storage_api::owner::{EcstoreDiskOption, EcstoreDiskStore, EcstoreEndpoint, ecstore_config_init, ecstore_new_disk}; @@ -61,6 +61,7 @@ pub use scanner::init_data_scanner; pub use scanner_io::{clear_dirty_usage_bucket, record_dirty_usage_bucket}; pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER}; use std::sync::atomic::{AtomicU64, Ordering}; +pub use storage_api::ScannerReplicationConfig as ReplicationConfig; static SCANNER_ACTIVE_WORK_UNITS: AtomicU64 = AtomicU64::new(0); static SCANNER_FOREGROUND_READ_ACTIVITY: AtomicU64 = AtomicU64::new(0); @@ -150,9 +151,9 @@ pub(crate) type Evaluator = EcstoreEvaluator; pub(crate) type Event = EcstoreEvent; pub(crate) type LcEventSrc = EcstoreLcEventSrc; pub(crate) type ObjectOpts = EcstoreObjectOpts; -pub(crate) type ReplicationConfig = EcstoreReplicationConfig; -pub(crate) type ReplicationHealQueueResult = EcstoreReplicationHealQueueResult; -pub(crate) type ReplicationQueueAdmission = EcstoreReplicationQueueAdmission; +pub(crate) type ReplicationHealObject = ScannerReplicationHealObject; +pub(crate) type ReplicationHealQueueResult = ScannerReplicationHealResult; +pub(crate) type ReplicationQueueAdmission = ScannerReplicationQueueAdmission; pub(crate) type ScanGuard = EcstoreScanGuard; pub(crate) type SetDisks = EcstoreSetDisks; pub(crate) type StorageError = EcstoreStorageError; @@ -310,7 +311,9 @@ pub(crate) async fn queue_replication_heal( rcfg: ReplicationConfig, retry_count: u32, ) -> ReplicationHealQueueResult { - EcstoreReplicationScannerBridge::queue_heal(bucket, oi, rcfg, retry_count).await + EcstoreReplicationScannerBridge::queue_heal(bucket, oi, rcfg.into_ecstore(), retry_count) + .await + .into() } pub(crate) fn resolve_scanner_object_store_handle() -> Option> { diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 20ff08542..8f8366b7a 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -41,9 +41,7 @@ use rustfs_common::metrics::{ IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, UpdateCurrentPathFn, current_path_updater, global_metrics, }; -use rustfs_filemeta::{ - MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ReplicateObjectInfo, ReplicationStatusType, ReplicationType, -}; +use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ReplicationStatusType}; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration}; use time::OffsetDateTime; @@ -54,10 +52,10 @@ use tracing::{debug, error, warn}; use crate::{ BucketVersioningSys, Disk, DiskError, DiskInfoOptions, Evaluator, Event, LcEventSrc, ListPathRawOptions, ObjectOpts, - ReplicationConfig, ReplicationQueueAdmission, ScannerDiskExt as _, ScannerLifecycleConfigExt as _, - ScannerReplicationConfigExt as _, ScannerVersioningConfigExt as _, StorageError, apply_expiry_rule, apply_transition_rule, - enqueue_runtime_newer_noncurrent, is_reserved_or_invalid_bucket, list_path_raw, path2_bucket_object, - path2_bucket_object_with_base_path, queue_replication_heal, scanner_is_erasure, + ReplicationConfig, ReplicationHealObject, ReplicationQueueAdmission, ScannerDiskExt as _, ScannerLifecycleConfigExt as _, + ScannerVersioningConfigExt as _, StorageError, apply_expiry_rule, apply_transition_rule, enqueue_runtime_newer_noncurrent, + is_reserved_or_invalid_bucket, list_path_raw, path2_bucket_object, path2_bucket_object_with_base_path, + queue_replication_heal, scanner_is_erasure, scanner_replication_config_for_lifecycle_eval, }; use crate::{ScannerObjectInfo as ObjectInfo, ScannerObjectToDelete as ObjectToDelete}; @@ -213,14 +211,14 @@ fn scanner_replication_work_update(admission: ReplicationQueueAdmission) -> Scan } } -fn scanner_replication_repair_kind(roi: &ReplicateObjectInfo) -> Option { - if roi.bucket.is_empty() && roi.name.is_empty() { +fn scanner_replication_repair_kind(roi: &ReplicationHealObject) -> Option { + if roi.is_empty_identity() { return None; } - if roi.op_type == ReplicationType::ExistingObject || roi.existing_obj_resync.must_resync() { + if roi.is_existing_object_repair() { Some(ScannerReplicationRepairKind::BucketExistingObject) - } else if !roi.version_purge_status.is_empty() { + } else if roi.has_version_purge_status() { Some(ScannerReplicationRepairKind::BucketVersionPurge) } else if roi.delete_marker { Some(ScannerReplicationRepairKind::BucketDeleteMarker) @@ -229,7 +227,7 @@ fn scanner_replication_repair_kind(roi: &ReplicateObjectInfo) -> Option, remotes: Option) -> Self { + Self(EcstoreReplicationConfig::new(config, remotes)) + } + + pub(crate) fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool { + !self.0.is_empty() + && self + .0 + .config + .as_ref() + .is_some_and(|config| EcstoreReplicationConfigurationExt::has_active_rules(config, prefix, recursive)) + } + + pub(crate) fn into_ecstore(self) -> EcstoreReplicationConfig { + self.0 + } +} + +pub(crate) fn scanner_replication_config_for_lifecycle_eval( + config: Option>, +) -> Option> { + config.map(|config| match Arc::try_unwrap(config) { + Ok(config) => Arc::new(config.into_ecstore()), + Err(config) => Arc::new(config.0.clone()), + }) +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub(crate) enum ScannerReplicationQueueAdmission { + #[default] + Skipped, + Queued, + Missed, +} + +impl From for ScannerReplicationQueueAdmission { + fn from(admission: EcstoreReplicationQueueAdmission) -> Self { + match admission { + EcstoreReplicationQueueAdmission::Skipped => Self::Skipped, + EcstoreReplicationQueueAdmission::Queued => Self::Queued, + EcstoreReplicationQueueAdmission::Missed => Self::Missed, + } + } +} + +#[derive(Debug, Clone, Default)] +pub(crate) struct ScannerReplicationHealObject { + pub(crate) bucket: String, + pub(crate) name: String, + pub(crate) size: i64, + pub(crate) delete_marker: bool, + pub(crate) target_statuses: HashMap, + version_purge_status: VersionPurgeStatusType, + existing_object: bool, + existing_object_resync: bool, +} + +impl ScannerReplicationHealObject { + #[cfg(test)] + pub(crate) fn new(bucket: impl Into, name: impl Into) -> Self { + Self { + bucket: bucket.into(), + name: name.into(), + ..Default::default() + } + } + + pub(crate) fn is_empty_identity(&self) -> bool { + self.bucket.is_empty() && self.name.is_empty() + } + + pub(crate) fn is_existing_object_repair(&self) -> bool { + self.existing_object || self.existing_object_resync + } + + pub(crate) fn has_version_purge_status(&self) -> bool { + !self.version_purge_status.is_empty() + } + + #[cfg(test)] + pub(crate) fn with_delete_marker(mut self) -> Self { + self.delete_marker = true; + self + } + + #[cfg(test)] + pub(crate) fn with_pending_version_purge(mut self) -> Self { + self.version_purge_status = VersionPurgeStatusType::Pending; + self + } + + #[cfg(test)] + pub(crate) fn with_existing_object(mut self) -> Self { + self.existing_object = true; + self + } + + #[cfg(test)] + pub(crate) fn with_existing_object_resync(mut self) -> Self { + self.existing_object_resync = true; + self + } +} + +impl From for ScannerReplicationHealObject { + fn from(object: ReplicateObjectInfo) -> Self { + Self { + bucket: object.bucket, + name: object.name, + size: object.size, + delete_marker: object.delete_marker, + target_statuses: object.target_statuses, + version_purge_status: object.version_purge_status, + existing_object: object.op_type == ReplicationType::ExistingObject, + existing_object_resync: object.existing_obj_resync.must_resync(), + } + } +} + +#[derive(Debug, Clone, Default)] +pub(crate) struct ScannerReplicationHealResult { + pub(crate) object_info: ScannerReplicationHealObject, + pub(crate) admission: ScannerReplicationQueueAdmission, +} + +impl From for ScannerReplicationHealResult { + fn from(result: EcstoreReplicationHealQueueResult) -> Self { + Self { + object_info: result.object_info.into(), + admission: result.admission.into(), + } + } +} + pub(crate) mod scan { pub(crate) use super::storage_contracts::{BucketOperations, BucketOptions, NamespaceLocking}; } diff --git a/docs/architecture/ecstore-api-facade-inventory.md b/docs/architecture/ecstore-api-facade-inventory.md index 4305d8968..db63d67ad 100644 --- a/docs/architecture/ecstore-api-facade-inventory.md +++ b/docs/architecture/ecstore-api-facade-inventory.md @@ -21,7 +21,7 @@ External `rustfs_ecstore::api` imports must stay in these local boundary files: | Boundary file | Current facade families | |---|---| | `rustfs/src/storage/storage_api.rs` | Broad RustFS storage owner bridge for admin, bucket, capacity, client, compression, cluster, config, data usage, disk, error, event, global bootstrap controls, runtime-source getters, layout, metrics, notification, rebalance, rio, rpc, set disk, storage, and tier. | -| `crates/scanner/src/storage_api.rs` | Scanner bridge for bucket lifecycle, replication, metadata, capacity, config, data usage, disk, error, runtime, set disk, storage, and tier. | +| `crates/scanner/src/storage_api.rs` | Scanner bridge for bucket lifecycle, replication, metadata, capacity, config, data usage, disk, error, runtime, set disk, storage, and tier. Replication queue config, admission, and heal object DTOs are projected into scanner-local types here. | | `crates/obs/src/metrics/storage_api.rs` | Metrics bridge for bucket bandwidth, lifecycle, replication, quota, capacity, data usage, error, runtime, and storage. | | `crates/iam/src/storage_api.rs` | IAM bridge for config, error, notification, runtime, and storage. | | `crates/heal/src/heal/storage_api.rs` | Heal bridge for data usage, disk, error, runtime, and storage. | diff --git a/docs/architecture/ecstore-module-split-plan.md b/docs/architecture/ecstore-module-split-plan.md index 56bf1c83c..4b7a9b0b5 100644 --- a/docs/architecture/ecstore-module-split-plan.md +++ b/docs/architecture/ecstore-module-split-plan.md @@ -196,6 +196,9 @@ Required contracts before crate movement: - `ReplicationScannerBridge`: scanner-originated replication heal scheduling is exposed through the contract type in `crates/ecstore/src/bucket/replication/replication_scanner_bridge.rs`. + Scanner consumers receive scanner-local replication config/admission/heal + object DTOs from `crates/scanner/src/storage_api.rs` instead of constructing + or inspecting replication queue DTOs directly. - `ReplicationTargetConfigBridge`: bucket target removal checks against replication target rules are exposed through the contract type in `crates/ecstore/src/bucket/replication/replication_target_config_bridge.rs`. diff --git a/scripts/check_architecture_migration_rules.sh b/scripts/check_architecture_migration_rules.sh index 91a6972f1..181a648bb 100755 --- a/scripts/check_architecture_migration_rules.sh +++ b/scripts/check_architecture_migration_rules.sh @@ -201,6 +201,7 @@ REPLICATION_FACADE_BYPASS_HITS_FILE="${TMP_DIR}/replication_facade_bypass_hits.t REPLICATION_FACADE_WILDCARD_EXPORT_HITS_FILE="${TMP_DIR}/replication_facade_wildcard_export_hits.txt" ADMIN_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/admin_replication_dto_boundary_bypass_hits.txt" APP_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/app_replication_dto_boundary_bypass_hits.txt" +SCANNER_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/scanner_replication_dto_boundary_bypass_hits.txt" OBS_REPLICATION_STATS_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/obs_replication_stats_boundary_bypass_hits.txt" REPLICATION_BANDWIDTH_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/replication_bandwidth_boundary_bypass_hits.txt" REPLICATION_CONFIG_STORE_BYPASS_HITS_FILE="${TMP_DIR}/replication_config_store_bypass_hits.txt" @@ -2593,6 +2594,24 @@ if [[ -s "$APP_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE" ]]; then report_failure "app replication ObjectOpts/MustReplicateOptions/bridge access must stay behind rustfs/src/app/storage_api.rs: $(paste -sd '; ' "$APP_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE")" fi +( + cd "$ROOT_DIR" + { + rg -n --with-filename 'rustfs_ecstore::api::bucket::replication' \ + crates/scanner/src crates/scanner/tests \ + --glob '*.rs' \ + --glob '!crates/scanner/src/storage_api.rs' || true + rg -n --with-filename '\b(ReplicateObjectInfo|ReplicationType|ResyncDecision|ResyncTargetDecision|EcstoreReplicationConfig|EcstoreReplicationHealQueueResult|EcstoreReplicationQueueAdmission)\b|ReplicationConfig\s*\{' \ + crates/scanner/src crates/scanner/tests \ + --glob '*.rs' \ + --glob '!crates/scanner/src/storage_api.rs' || true + } +) >"$SCANNER_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE" + +if [[ -s "$SCANNER_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE" ]]; then + report_failure "scanner replication queue DTO access must stay behind crates/scanner/src/storage_api.rs: $(paste -sd '; ' "$SCANNER_REPLICATION_DTO_BOUNDARY_BYPASS_HITS_FILE")" +fi + ( cd "$ROOT_DIR" {