From a96dd7d28921e5134e13da1ad43727673cc70735 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 26 Aug 2026 21:13:18 +0800 Subject: [PATCH] refactor: migrate consumers off rustfs-common heal/scanner shims (#6623) * refactor(ecstore): import heal/scanner contracts crates directly (backlog#1843) * refactor(heal): import heal/scanner contracts crates directly (backlog#1843) * refactor(lifecycle): import heal/scanner contracts crates directly (backlog#1843) * refactor(obs): import heal/scanner contracts crates directly (backlog#1843) * refactor(protos): import heal/scanner contracts crates directly (backlog#1843) * refactor(scanner): import heal/scanner contracts crates directly (backlog#1843) * refactor(rustfs): import heal/scanner contracts crates directly (backlog#1843) --- Cargo.lock | 11 ++++ crates/ecstore/Cargo.toml | 2 + .../lifecycle/bucket_lifecycle_audit.rs | 2 +- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 8 +-- .../src/bucket/lifecycle/replication_sink.rs | 4 +- crates/ecstore/src/bucket/metadata_sys.rs | 2 +- crates/ecstore/src/bucket/quota/checker.rs | 10 ++-- .../ecstore/src/cluster/rpc/peer_s3_client.rs | 2 +- crates/ecstore/src/config/com.rs | 4 +- crates/ecstore/src/core/pools.rs | 2 +- crates/ecstore/src/core/sets.rs | 4 +- .../src/diagnostics/admin_server_info.rs | 2 +- .../src/ecstore_validation_blackbox.rs | 4 +- crates/ecstore/src/layout/set_heal.rs | 2 +- crates/ecstore/src/object_api/types.rs | 2 +- .../ecstore/src/services/metrics_realtime.rs | 53 ++++++++++--------- .../ecstore/src/services/notification_sys.rs | 2 +- .../src/set_disk/core/io_primitives.rs | 37 +++++++------ crates/ecstore/src/set_disk/mod.rs | 22 ++++---- crates/ecstore/src/set_disk/ops/heal.rs | 2 +- crates/ecstore/src/set_disk/ops/multipart.rs | 6 +-- crates/ecstore/src/set_disk/ops/object.rs | 16 +++--- crates/ecstore/src/set_disk/read.rs | 14 +++-- crates/ecstore/src/store/heal.rs | 2 +- crates/ecstore/src/store/init.rs | 4 +- crates/ecstore/src/store/mod.rs | 2 +- crates/ecstore/src/store/rebalance.rs | 4 +- .../tests/ecstore_contract_compat_test.rs | 2 +- crates/heal/Cargo.toml | 1 + crates/heal/src/heal/channel.rs | 10 ++-- crates/heal/src/heal/erasure_healer.rs | 6 +-- crates/heal/src/heal/manager.rs | 4 +- crates/heal/src/heal/manager/tests.rs | 2 +- crates/heal/src/heal/mrf_queue.rs | 4 +- crates/heal/src/heal/storage.rs | 2 +- crates/heal/src/heal/task.rs | 2 +- crates/heal/src/lib.rs | 4 +- .../heal_b5_versioned_regression_test.rs | 2 +- .../tests/heal_b920_subquorum_union_test.rs | 2 +- crates/heal/tests/heal_bug_fixes_test.rs | 8 +-- crates/heal/tests/heal_integration_test.rs | 2 +- crates/lifecycle/Cargo.toml | 1 + crates/lifecycle/src/core.rs | 2 +- crates/lifecycle/src/evaluator.rs | 4 +- crates/lifecycle/src/lib.rs | 2 +- crates/obs/Cargo.toml | 2 + crates/obs/src/metrics/collectors/scanner.rs | 2 +- crates/obs/src/metrics/stats_collector.rs | 16 +++--- crates/protos/Cargo.toml | 1 + crates/protos/src/heal_control.rs | 4 +- crates/scanner/Cargo.toml | 2 + crates/scanner/src/data_usage_define.rs | 2 +- crates/scanner/src/lib.rs | 2 +- crates/scanner/src/remote_scanner/stream.rs | 4 +- crates/scanner/src/scanner.rs | 12 ++--- crates/scanner/src/scanner_folder.rs | 16 +++--- crates/scanner/src/scanner_folder/tests.rs | 8 +-- crates/scanner/src/scanner_io.rs | 6 ++- crates/scanner/src/scanner_io/guards.rs | 6 +-- crates/scanner/src/scanner_io/io_disk.rs | 4 +- crates/scanner/src/sleeper.rs | 2 +- rustfs/Cargo.toml | 2 + rustfs/src/admin/handlers/heal.rs | 44 ++++++++------- rustfs/src/admin/handlers/scanner.rs | 4 +- rustfs/src/admin/storage_api.rs | 2 +- rustfs/src/app/capacity_dirty_scope_test.rs | 2 +- rustfs/src/app/storage_api.rs | 4 +- rustfs/src/cluster_snapshot.rs | 2 +- rustfs/src/startup_runtime_sources.rs | 2 +- rustfs/src/storage/rpc/node_service.rs | 25 +++++---- rustfs/src/storage/rpc/node_service/heal.rs | 4 +- rustfs/src/storage/storage_api.rs | 4 +- 72 files changed, 263 insertions(+), 206 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 45fef954d..7b83156ba 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9370,6 +9370,7 @@ dependencies = [ "rustfs-extension-schema", "rustfs-filemeta", "rustfs-heal", + "rustfs-heal-contracts", "rustfs-iam", "rustfs-io-core", "rustfs-io-metrics", @@ -9393,6 +9394,7 @@ dependencies = [ "rustfs-s3select-api", "rustfs-s3select-query", "rustfs-scanner", + "rustfs-scanner-contracts", "rustfs-security-governance", "rustfs-signer", "rustfs-storage-api", @@ -9632,6 +9634,7 @@ dependencies = [ "rustfs-data-usage", "rustfs-erasure-codec", "rustfs-filemeta", + "rustfs-heal-contracts", "rustfs-io-metrics", "rustfs-lifecycle", "rustfs-lock", @@ -9645,6 +9648,7 @@ dependencies = [ "rustfs-rio-v2", "rustfs-s3-client", "rustfs-s3-types", + "rustfs-scanner-contracts", "rustfs-signer", "rustfs-storage-api", "rustfs-tls-runtime", @@ -9755,6 +9759,7 @@ dependencies = [ "rustfs-concurrency", "rustfs-config", "rustfs-ecstore", + "rustfs-heal-contracts", "rustfs-lock", "rustfs-madmin", "rustfs-storage-api", @@ -9996,6 +10001,7 @@ dependencies = [ "rustfs-common", "rustfs-config", "rustfs-replication", + "rustfs-scanner-contracts", "rustfs-storage-api", "s3s", "serial_test", @@ -10190,9 +10196,11 @@ dependencies = [ "rustfs-common", "rustfs-config", "rustfs-ecstore", + "rustfs-heal-contracts", "rustfs-iam", "rustfs-io-metrics", "rustfs-notify", + "rustfs-scanner-contracts", "rustfs-security-governance", "rustfs-storage-api", "rustfs-utils", @@ -10317,6 +10325,7 @@ dependencies = [ "rmp-serde", "rustfs-common", "rustfs-config", + "rustfs-heal-contracts", "rustfs-io-metrics", "rustfs-tls-runtime", "rustfs-utils", @@ -10544,8 +10553,10 @@ dependencies = [ "rustfs-data-usage", "rustfs-ecstore", "rustfs-filemeta", + "rustfs-heal-contracts", "rustfs-lock", "rustfs-s3-types", + "rustfs-scanner-contracts", "rustfs-storage-api", "rustfs-utils", "s3s", diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index fb223489c..fc488fd77 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -143,6 +143,8 @@ rustfs-config = { workspace = true, features = ["notify", "audit", "server-confi rustfs-concurrency.workspace = true rustfs-credentials = { workspace = true } rustfs-common.workspace = true +rustfs-heal-contracts.workspace = true +rustfs-scanner-contracts.workspace = true rustfs-policy.workspace = true rustfs-protos.workspace = true rustfs-replication.workspace = true diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs index cfeb644a8..06ea063a7 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs @@ -16,8 +16,8 @@ use super::runtime_boundary as runtime_sources; use crate::bucket::lifecycle::lifecycle; use crate::object_api::ObjectInfo; use crate::services::event_notification::{EventArgs, send_event}; -use rustfs_common::metrics::IlmAction; use rustfs_s3_types::EventName; +use rustfs_scanner_contracts::metrics::IlmAction; const LIFECYCLE_EXPIRY_USER_AGENT: &str = "Internal: [ILM-Expiry]"; const LIFECYCLE_TRANSITION_USER_AGENT: &str = "Internal: [ILM-Transition]"; diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 5e13ba918..c3c1f67f0 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -72,9 +72,6 @@ use crate::store::ECStore; use async_channel::{Receiver as A_Receiver, Sender as A_Sender, bounded}; use http::HeaderMap; use rand::RngExt as _; -use rustfs_common::metrics::{ - IlmAction, Metrics, ScannerLifecycleExpiryStateUpdate, ScannerLifecycleTransitionStateUpdate, global_metrics, -}; use rustfs_config::{ DEFAULT_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_QUEUE_SEND_TIMEOUT_MS, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX, DEFAULT_TRANSITION_WORKERS_CAP, ENV_MAX_EXPIRY_WORKERS, ENV_TRANSITION_QUEUE_CAPACITY, ENV_TRANSITION_QUEUE_SEND_TIMEOUT_MS, @@ -84,6 +81,9 @@ use rustfs_data_usage::TierStats; use rustfs_filemeta::{ FileInfo, FileInfoOpts, NULL_VERSION_ID, RestoreStatusOps, TRANSITION_COMPLETE, get_file_info, is_restored_object_on_disk, }; +use rustfs_scanner_contracts::metrics::{ + IlmAction, Metrics, ScannerLifecycleExpiryStateUpdate, ScannerLifecycleTransitionStateUpdate, global_metrics, +}; use rustfs_utils::{ get_env_i64, get_env_usize, path::encode_dir_object, @@ -5460,11 +5460,11 @@ mod tests { use futures::FutureExt; #[cfg(feature = "test-util")] use http::HeaderMap; - use rustfs_common::metrics::{IlmAction, global_metrics}; use rustfs_config::ENV_MAX_EXPIRY_WORKERS; use rustfs_config::ENV_TRANSITION_WORKERS_ABSOLUTE_MAX; use rustfs_data_usage::TierStats; use rustfs_filemeta::{FileInfo, FileMeta}; + use rustfs_scanner_contracts::metrics::{IlmAction, global_metrics}; use s3s::dto::{ BucketLifecycleConfiguration, DefaultRetention, ExpirationStatus, LifecycleExpiration, LifecycleRule, MetadataEntry, ObjectLockConfiguration, ObjectLockEnabled, ObjectLockRetentionMode, ObjectLockRule, OutputLocation, RestoreRequest, diff --git a/crates/ecstore/src/bucket/lifecycle/replication_sink.rs b/crates/ecstore/src/bucket/lifecycle/replication_sink.rs index 32cd6b57a..70b700883 100644 --- a/crates/ecstore/src/bucket/lifecycle/replication_sink.rs +++ b/crates/ecstore/src/bucket/lifecycle/replication_sink.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use rustfs_common::metrics::IlmAction; +use rustfs_scanner_contracts::metrics::IlmAction; use crate::bucket::lifecycle::lifecycle::ObjectOpts; use crate::bucket::replication::ReplicationLifecycleBridge; @@ -77,7 +77,7 @@ mod tests { use crate::bucket::replication::{DeleteReplicationConfigSnapshot, ReplicationObjectBridge}; use crate::object_api::{ObjectInfo, ObjectOptions}; use crate::storage_api_contracts::object::ObjectToDelete; - use rustfs_common::metrics::IlmAction; + use rustfs_scanner_contracts::metrics::IlmAction; use s3s::dto::{ BucketVersioningStatus, DeleteMarkerReplication, DeleteMarkerReplicationStatus, DeleteReplication, DeleteReplicationStatus, Destination, ReplicationConfiguration, ReplicationRule, ReplicationRuleStatus, diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 2cbb0eff4..ba166d7f1 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -27,7 +27,7 @@ use crate::storage_api_contracts::heal::HealOperations as _; use crate::storage_api_contracts::namespace::NamespaceLocking as _; use crate::store::{ECStore, await_bucket_namespace_operation}; use futures::future::join_all; -use rustfs_common::heal_channel::HealOpts; +use rustfs_heal_contracts::heal_channel::HealOpts; use rustfs_policy::policy::BucketPolicy; use s3s::dto::ReplicationConfiguration; use s3s::dto::{ diff --git a/crates/ecstore/src/bucket/quota/checker.rs b/crates/ecstore/src/bucket/quota/checker.rs index 0d71095e3..44a113652 100644 --- a/crates/ecstore/src/bucket/quota/checker.rs +++ b/crates/ecstore/src/bucket/quota/checker.rs @@ -15,8 +15,8 @@ use super::{BucketQuota, QuotaCheckResult, QuotaError, QuotaOperation}; use crate::bucket::metadata_sys::{BucketMetadataSys, update, update_if_incarnation}; use crate::data_usage::get_bucket_usage_memory; -use rustfs_common::metrics::Metric; use rustfs_config::QUOTA_CONFIG_FILE; +use rustfs_scanner_contracts::metrics::Metric; use std::sync::Arc; use std::time::Instant; use time::OffsetDateTime; @@ -120,9 +120,9 @@ impl QuotaChecker { let duration = start_time.elapsed(); // inc_time is now a plain fn (not async) — no .await needed. - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaCheck, duration); + rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaCheck, duration); if !allowed { - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaViolation, duration); + rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaViolation, duration); } Ok(result) @@ -185,7 +185,7 @@ impl QuotaChecker { .await .map_err(QuotaError::StorageError)?; - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed()); + rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed()); Ok(updated_at) } @@ -206,7 +206,7 @@ impl QuotaChecker { } .map_err(QuotaError::StorageError)?; - rustfs_common::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed()); + rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed()); Ok(updated_at) } diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index cbd944fb0..8e5ebcecc 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -36,7 +36,7 @@ use crate::{ }; use async_trait::async_trait; use futures::future::join_all; -use rustfs_common::heal_channel::{DriveState, HealItemType, HealOpts, RUSTFS_RESERVED_BUCKET}; +use rustfs_heal_contracts::heal_channel::{DriveState, HealItemType, HealOpts, RUSTFS_RESERVED_BUCKET}; use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem}; use rustfs_protos::proto_gen::node_service::node_service_client::NodeServiceClient; use rustfs_protos::proto_gen::node_service::{ diff --git a/crates/ecstore/src/config/com.rs b/crates/ecstore/src/config/com.rs index 98409f513..16f9f4c20 100644 --- a/crates/ecstore/src/config/com.rs +++ b/crates/ecstore/src/config/com.rs @@ -27,7 +27,6 @@ use crate::storage_api_contracts::{ }; use crate::store::ECStore; use http::HeaderMap; -use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_config::audit::{ AUDIT_AMQP_KEYS, AUDIT_AMQP_SUB_SYS, AUDIT_KAFKA_KEYS, AUDIT_KAFKA_SUB_SYS, AUDIT_MQTT_KEYS, AUDIT_MQTT_SUB_SYS, AUDIT_MYSQL_KEYS, AUDIT_MYSQL_SUB_SYS, AUDIT_NATS_KEYS, AUDIT_NATS_SUB_SYS, AUDIT_POSTGRES_KEYS, AUDIT_POSTGRES_SUB_SYS, @@ -46,6 +45,7 @@ use rustfs_config::{ SCANNER_SUB_SYS, }; use rustfs_filemeta::FileInfo; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; use serde_json::{Map, Value}; use std::collections::{HashMap, HashSet}; use std::sync::LazyLock; @@ -4710,7 +4710,7 @@ mod tests { decode_persisted_server_config, fallback_server_config_after_corruption, is_server_config_corrupt_error, read_config_without_migrate_with_recovery, replace_server_config_decrypt_fn_for_test, }; - use rustfs_common::heal_channel::HealOpts; + use rustfs_heal_contracts::heal_channel::HealOpts; use std::sync::Mutex; /// Bytes mirroring issue #4156: a bitrot-corrupted `config.json` whose diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 0f188db18..6570db998 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -70,8 +70,8 @@ use http::HeaderMap; #[cfg(test)] use rmp_serde::Deserializer; use rmp_serde::Serializer; -use rustfs_common::heal_channel::HealOpts; use rustfs_filemeta::{FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams}; +use rustfs_heal_contracts::heal_channel::HealOpts; use rustfs_utils::crypto::{hex_sha256, is_sha256_checksum}; use rustfs_utils::path::{encode_dir_object, path_join, path_to_bucket_object, path_to_bucket_object_with_base_path}; use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration, ReplicationConfiguration}; diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 648b7cf0d..0a814865f 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -46,9 +46,9 @@ use futures::{ stream::{FuturesUnordered, StreamExt}, }; use http::HeaderMap; -use rustfs_common::heal_channel::HealOpts; -use rustfs_common::heal_channel::{DriveState, HealItemType}; use rustfs_filemeta::FileInfo; +use rustfs_heal_contracts::heal_channel::HealOpts; +use rustfs_heal_contracts::heal_channel::{DriveState, HealItemType}; use rustfs_lock::NamespaceLockWrapper; use rustfs_madmin::heal_commands::HealResultItem; use rustfs_utils::{crc_hash, path::path_join_buf, sip_hash}; diff --git a/crates/ecstore/src/diagnostics/admin_server_info.rs b/crates/ecstore/src/diagnostics/admin_server_info.rs index 017b462a2..3bbf7cf05 100644 --- a/crates/ecstore/src/diagnostics/admin_server_info.rs +++ b/crates/ecstore/src/diagnostics/admin_server_info.rs @@ -26,7 +26,7 @@ use crate::{ use crate::data_usage::load_data_usage_cache; use crate::storage_api_contracts::admin::StorageAdminApi; use crate::storage_api_contracts::bucket::BucketOptions; -use rustfs_common::heal_channel::DriveState; +use rustfs_heal_contracts::heal_channel::DriveState; use rustfs_madmin::{ BackendDisks, Disk, ErasureSetInfo, ITEM_INITIALIZING, ITEM_OFFLINE, ITEM_ONLINE, ITEM_UNKNOWN, InfoMessage, MemStats, ServerProperties, diff --git a/crates/ecstore/src/ecstore_validation_blackbox.rs b/crates/ecstore/src/ecstore_validation_blackbox.rs index 30e309e4f..4db7e5104 100644 --- a/crates/ecstore/src/ecstore_validation_blackbox.rs +++ b/crates/ecstore/src/ecstore_validation_blackbox.rs @@ -208,7 +208,7 @@ async fn blackbox_get_restores_body_after_one_shard_file_is_removed() { // Serialized: forces the reader-setup strategy through a process-global env var. #[serial_test::serial] async fn blackbox_heal_requests_preserve_repair_scope() { - use rustfs_common::heal_channel::{ + use rustfs_heal_contracts::heal_channel::{ HealAdmissionResult, HealChannelCommand, HealChannelPriority, HealChannelReceiver, HealChannelRequest, HealRequestSource, }; @@ -238,7 +238,7 @@ async fn blackbox_heal_requests_preserve_repair_scope() { // (failing their submitter, which releases their dedup reservation), and // fail fast once the receiver drops at test end. Tests that must observe a // deterministic channel state serialize under the same serial key. - let mut heal_rx = rustfs_common::heal_channel::init_heal_channel() + let mut heal_rx = rustfs_heal_contracts::heal_channel::init_heal_channel() .expect("this must be the only ecstore test that owns the heal channel receiver"); // Ordinary PUTs use the same admission channel as read repair. A single diff --git a/crates/ecstore/src/layout/set_heal.rs b/crates/ecstore/src/layout/set_heal.rs index 702016f24..e32beb84f 100644 --- a/crates/ecstore/src/layout/set_heal.rs +++ b/crates/ecstore/src/layout/set_heal.rs @@ -14,7 +14,7 @@ use crate::disk::{DiskInfo, error::DiskError}; use crate::layout::{endpoints::Endpoints, format::FormatV3}; -use rustfs_common::heal_channel::DriveState; +use rustfs_heal_contracts::heal_channel::DriveState; use rustfs_madmin::heal_commands::HealDriveInfo; pub(crate) fn formats_to_drives_info( diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index d91d96a0d..280dd10cc 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -281,7 +281,7 @@ pub struct QuotaAdmission { pub struct LifecycleDeleteAllRequest { pub(crate) version_id: Option, pub(crate) delete_marker: bool, - pub(crate) action: rustfs_common::metrics::IlmAction, + pub(crate) action: rustfs_scanner_contracts::metrics::IlmAction, pub(crate) rule_id: String, pub(crate) phase: LifecycleDeleteAllPhase, } diff --git a/crates/ecstore/src/services/metrics_realtime.rs b/crates/ecstore/src/services/metrics_realtime.rs index 290504fcb..c6056414d 100644 --- a/crates/ecstore/src/services/metrics_realtime.rs +++ b/crates/ecstore/src/services/metrics_realtime.rs @@ -18,7 +18,7 @@ use crate::storage_api_contracts::admin::StorageAdminApi; #[cfg(test)] use chrono::Utc; use jiff::Timestamp; -use rustfs_common::{heal_channel::DriveState, metrics::global_metrics}; +use rustfs_heal_contracts::heal_channel::DriveState; use rustfs_io_metrics::internode_metrics::global_internode_metrics; use rustfs_madmin::metrics::{ DiskIOStats, DiskMetric, LastMinute as MadminLastMinute, NetDevLine, NetMetrics, RPCMetrics, RealtimeMetrics, @@ -32,6 +32,7 @@ use rustfs_madmin::metrics::{ ScannerSourceCycleSnapshot as MadminScannerSourceCycleSnapshot, ScannerSourceWorkSnapshot as MadminScannerSourceWorkSnapshot, ScannerUsageFreshnessSnapshot as MadminScannerUsageFreshnessSnapshot, TimedAction as MadminTimedAction, }; +use rustfs_scanner_contracts::metrics::global_metrics; use rustfs_utils::os::get_drive_stats; use serde::{Deserialize, Serialize}; use std::collections::{HashMap, HashSet}; @@ -81,7 +82,7 @@ fn unix_millis_to_jiff_timestamp(millis: u64, fallback: Timestamp) -> Timestamp } } -fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsReport) -> MadminScannerMetrics { +fn to_madmin_scanner_metrics(metrics: rustfs_scanner_contracts::metrics::ScannerMetricsReport) -> MadminScannerMetrics { MadminScannerMetrics { collected_at: metrics.collected_at, current_cycle: metrics.current_cycle, @@ -563,8 +564,8 @@ async fn collect_local_disks_metrics(disks: &HashSet) -> HashMap Option { disks.push(rustfs_madmin::Disk { endpoint: ep.to_string(), state: if online { - rustfs_common::heal_channel::DriveState::Ok.to_string() + rustfs_heal_contracts::heal_channel::DriveState::Ok.to_string() } else { ItemState::Offline.to_string().to_owned() }, diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 2cc3935a6..6d2c6cf6c 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -1364,7 +1364,7 @@ pub(in crate::set_disk) enum ReadRepairAdmissionOutcome { pub(in crate::set_disk) type ReadRepairAdmissionFuture = Pin + Send>>; pub(in crate::set_disk) type ReadRepairAdmissionSubmitter = - fn(rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture; + fn(rustfs_heal_contracts::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture; pub(in crate::set_disk) struct ReadRepairHealSubmission<'a> { pub(in crate::set_disk) bucket: &'a str, @@ -1385,7 +1385,7 @@ pub(in crate::set_disk) struct ReadRepairHealSubmission<'a> { } pub(in crate::set_disk) fn send_read_repair_heal_request( - request: rustfs_common::heal_channel::HealChannelRequest, + request: rustfs_heal_contracts::heal_channel::HealChannelRequest, ) -> ReadRepairAdmissionFuture { Box::pin(async { match send_heal_request_with_admission(request).await { @@ -1453,7 +1453,7 @@ pub(in crate::set_disk) async fn submit_read_repair_heal_with_submitter( let _ = rustfs_common::mrf_channel::try_send_mrf_intent_typed(kind, bucket, object, version_uuid, Some(scope)); } - let mut request = rustfs_common::heal_channel::create_heal_request_with_options( + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( bucket.to_string(), Some(object.to_string()), false, @@ -3580,7 +3580,7 @@ pub(in crate::set_disk) async fn finish_rename_tail_heal< tail_drain: tokio::task::JoinHandle>, guard_release: tokio::sync::oneshot::Receiver, guards: Guards, - request: rustfs_common::heal_channel::HealChannelRequest, + request: rustfs_heal_contracts::heal_channel::HealChannelRequest, finalize: Finalize, cleanup: Cleanup, submit: Submit, @@ -3590,7 +3590,7 @@ pub(in crate::set_disk) async fn finish_rename_tail_heal< FinalizeFuture: Future + Send, Cleanup: FnOnce(Guards, Vec) -> CleanupFuture + Send, CleanupFuture: Future + Send, - Submit: FnOnce(rustfs_common::heal_channel::HealChannelRequest) -> SubmitFuture + Send, + Submit: FnOnce(rustfs_heal_contracts::heal_channel::HealChannelRequest) -> SubmitFuture + Send, SubmitFuture: Future + Send, { let (needs_heal, tail_cleanup, tail_complete) = match tail_drain.await { @@ -4939,16 +4939,17 @@ impl SetDisks { // reclaim_orphan_data_dirs. Reuses the existing heal channel, which // deduplicates and back-pressures via admission; failures only drop // the return value (same shape as multipart's existing heal enqueue). - let _ = - rustfs_common::heal_channel::send_heal_request(rustfs_common::heal_channel::create_heal_request_with_options( + let _ = rustfs_heal_contracts::heal_channel::send_heal_request( + rustfs_heal_contracts::heal_channel::create_heal_request_with_options( bucket.to_string(), Some(object.to_string()), false, - Some(rustfs_common::heal_channel::HealChannelPriority::Normal), + Some(rustfs_heal_contracts::heal_channel::HealChannelPriority::Normal), Some(self.pool_index), Some(self.set_index), - )) - .await; + ), + ) + .await; } } @@ -6800,18 +6801,24 @@ mod tests { write_raw_file_meta_unchecked(disk, bucket, object, metadata).await; } - fn failed_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + fn failed_read_repair_submitter( + _request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) -> ReadRepairAdmissionFuture { Box::pin(async { ReadRepairAdmissionOutcome::Failed("injected submit failure".to_string()) }) } - fn accepted_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + fn accepted_read_repair_submitter( + _request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) -> ReadRepairAdmissionFuture { Box::pin(async { ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Accepted) }) } - fn dropped_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + fn dropped_read_repair_submitter( + _request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) -> ReadRepairAdmissionFuture { Box::pin(async { ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Dropped( - rustfs_common::heal_channel::HealAdmissionDropReason::PolicyDropped, + rustfs_heal_contracts::heal_channel::HealAdmissionDropReason::PolicyDropped, )) }) } @@ -8853,7 +8860,7 @@ mod tests { tail_drain, released, (), - rustfs_common::heal_channel::HealChannelRequest::default(), + rustfs_heal_contracts::heal_channel::HealChannelRequest::default(), move || async move { *finalize_captured.lock().expect("finalize recorder should not poison") = true; }, diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index e26fcdd40..0d0f1fb57 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -130,15 +130,15 @@ use http::HeaderMap; use md5::{Digest as Md5Digest, Md5}; use rand::{Rng, seq::SliceRandom}; use regex::Regex; -use rustfs_common::heal_channel::{ - DriveState, HealAdmissionResult, HealChannelPriority, HealItemType, HealOpts, HealRequestSource, HealScanMode, - send_heal_disk, send_heal_request_with_admission, -}; use rustfs_config::MI_B; use rustfs_filemeta::{ FileInfo, FileMeta, FileMetaShallowVersion, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo, RawFileInfo, file_info_from_raw, merge_file_meta_versions, }; +use rustfs_heal_contracts::heal_channel::{ + DriveState, HealAdmissionResult, HealChannelPriority, HealItemType, HealOpts, HealRequestSource, HealScanMode, + send_heal_disk, send_heal_request_with_admission, +}; use rustfs_io_metrics::{ record_object_lock_diag_acquire_duration, record_object_lock_diag_enabled, record_object_lock_diag_hold_duration, record_object_lock_diag_slow_acquire, record_object_lock_diag_slow_hold, @@ -3111,8 +3111,9 @@ pub struct SetDisks { #[cfg(test)] storage_class_config_override: Arc>>>, #[cfg(test)] - rename_tail_heal_capture: - Arc>>>, + rename_tail_heal_capture: Arc< + std::sync::Mutex>>, + >, } // DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones @@ -3388,7 +3389,10 @@ impl DiskHealthEntry { } impl SetDisks { - pub(in crate::set_disk) async fn submit_rename_tail_heal(&self, request: rustfs_common::heal_channel::HealChannelRequest) { + pub(in crate::set_disk) async fn submit_rename_tail_heal( + &self, + request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) { #[cfg(test)] { let capture = self @@ -3402,13 +3406,13 @@ impl SetDisks { } } - let _ = rustfs_common::heal_channel::send_heal_request(request).await; + let _ = rustfs_heal_contracts::heal_channel::send_heal_request(request).await; } #[cfg(test)] pub(in crate::set_disk) fn capture_test_rename_tail_heals( &self, - ) -> tokio::sync::mpsc::UnboundedReceiver { + ) -> tokio::sync::mpsc::UnboundedReceiver { let (capture, requests) = tokio::sync::mpsc::unbounded_channel(); let mut slot = self .rename_tail_heal_capture diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index adcea3a60..d1a5c9281 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -2386,8 +2386,8 @@ mod heal_result_report_tests { store::init_format::{load_format_erasure, save_format_file}, }; use bytes::Bytes; - use rustfs_common::heal_channel::{DriveState, HealOpts, HealScanMode}; use rustfs_filemeta::{BLOCK_SIZE_V2, FileInfo, ObjectPartInfo, TRANSITION_COMPLETE}; + use rustfs_heal_contracts::heal_channel::{DriveState, HealOpts, HealScanMode}; use std::sync::{Arc, Mutex}; use tempfile::TempDir; use time::OffsetDateTime; diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index cd9b83623..7aafc24d2 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -2749,7 +2749,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { needs_immediate_heal = rename_commit.needs_immediate_heal(); if let Some(rename_tail_drain) = rename_commit.tail_drain.take() { tail_owns_staging_cleanup = true; - let mut request = rustfs_common::heal_channel::create_heal_request_with_options( + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( commit_bucket.clone(), Some(commit_object.clone()), false, @@ -2846,7 +2846,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let committed_file_info = rename_commit.committed_file_info; if needs_immediate_heal { - let mut request = rustfs_common::heal_channel::create_heal_request_with_options( + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( commit_bucket.clone(), Some(commit_object.clone()), false, @@ -2859,7 +2859,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .or_else(|| commit_version_suspended.then(Uuid::nil)) .map(|version_id| version_id.to_string()); tokio::spawn(async move { - let _ = rustfs_common::heal_channel::send_heal_request(request).await; + let _ = rustfs_heal_contracts::heal_channel::send_heal_request(request).await; }); } diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 87f46e868..8f550be82 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -352,7 +352,7 @@ mod lifecycle_delete_all_plan_tests { crate::object_api::LifecycleDeleteAllRequest { version_id: Some(version_id), delete_marker: true, - action: rustfs_common::metrics::IlmAction::DelMarkerDeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DelMarkerDeleteAllVersionsAction, rule_id: "rule".to_string(), phase: crate::object_api::LifecycleDeleteAllPhase::Preflight, } @@ -545,7 +545,7 @@ mod lifecycle_delete_all_plan_tests { let request = crate::object_api::LifecycleDeleteAllRequest { version_id: None, delete_marker: false, - action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction, rule_id: "rule".to_string(), phase: crate::object_api::LifecycleDeleteAllPhase::Preflight, }; @@ -3242,7 +3242,7 @@ impl SetDisks { needs_immediate_heal = rename_commit.needs_immediate_heal(); if let Some(rename_tail_drain) = rename_commit.tail_drain.take() { tail_owns_tmp_cleanup = true; - let mut request = rustfs_common::heal_channel::create_heal_request_with_options( + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( commit_bucket.clone(), Some(commit_object.clone()), false, @@ -3350,7 +3350,7 @@ impl SetDisks { let mut fi = rename_commit.committed_file_info; if needs_immediate_heal { - let mut request = rustfs_common::heal_channel::create_heal_request_with_options( + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( commit_bucket.clone(), Some(commit_object.clone()), false, @@ -3362,7 +3362,7 @@ impl SetDisks { .or_else(|| commit_version_suspended.then(Uuid::nil)) .map(|version_id| version_id.to_string()); tokio::spawn(async move { - let _ = rustfs_common::heal_channel::send_heal_request(request).await; + let _ = rustfs_heal_contracts::heal_channel::send_heal_request(request).await; }); } @@ -7076,7 +7076,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { Some(scope), ); } - let mut request = rustfs_common::heal_channel::create_heal_request_with_options( + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( bucket.to_string(), Some(object.to_string()), false, @@ -7085,7 +7085,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { Some(self.set_index), ); request.object_version_id = (!version_id.is_empty()).then(|| version_id.to_string()); - if let Err(e) = rustfs_common::heal_channel::send_heal_request(request).await { + if let Err(e) = rustfs_heal_contracts::heal_channel::send_heal_request(request).await { warn!( bucket, object, @@ -16546,7 +16546,7 @@ mod delete_objects_lock_gating_tests { lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest { version_id: Some(trigger_version_id), delete_marker: false, - action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction, rule_id: "rule".to_string(), phase: crate::object_api::LifecycleDeleteAllPhase::History, }), diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 15e7e4e2a..569ba57c0 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -1896,7 +1896,7 @@ fn is_get_object_metadata_cache_request_eligible(bucket: &str, opts: &ObjectOpti #[cfg(test)] mod metadata_cache_tests { use super::*; - use rustfs_common::heal_channel::HealAdmissionDropReason; + use rustfs_heal_contracts::heal_channel::HealAdmissionDropReason; use serial_test::serial; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Mutex, OnceLock}; @@ -1974,7 +1974,9 @@ mod metadata_cache_tests { } } - fn slow_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + fn slow_read_repair_submitter( + _request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) -> ReadRepairAdmissionFuture { SLOW_READ_REPAIR_SUBMITTER_CALLS.fetch_add(1, Ordering::Relaxed); Box::pin(async { tokio::time::sleep(Duration::from_millis(250)).await; @@ -1982,14 +1984,18 @@ mod metadata_cache_tests { }) } - fn dropped_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + fn dropped_read_repair_submitter( + _request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) -> ReadRepairAdmissionFuture { DROPPED_READ_REPAIR_SUBMITTER_CALLS.fetch_add(1, Ordering::Relaxed); Box::pin(async { ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped)) }) } - fn capture_read_repair_submitter(request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + fn capture_read_repair_submitter( + request: rustfs_heal_contracts::heal_channel::HealChannelRequest, + ) -> ReadRepairAdmissionFuture { CAPTURED_READ_REPAIR_CALLS.fetch_add(1, Ordering::Relaxed); *CAPTURED_READ_REPAIR_PRIORITY.lock().expect("capture mutex poisoned") = Some(request.priority); Box::pin(async { diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 206218492..9159d38b3 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -608,7 +608,7 @@ mod tests { use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations}; use crate::store::init_format::{load_format_erasure, save_format_file}; use crate::store::init_local_disks_with_instance_ctx; - use rustfs_common::heal_channel::DriveState; + use rustfs_heal_contracts::heal_channel::DriveState; use tokio_util::sync::CancellationToken; #[derive(Debug)] diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 65292980a..21a9bb454 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -9537,7 +9537,7 @@ mod tests { .expect("unknown transition metadata should be written"); } let lifecycle_event = crate::bucket::lifecycle::lifecycle::Event { - action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction, rule_id: "delete-all-versions".to_string(), ..Default::default() }; @@ -9709,7 +9709,7 @@ mod tests { lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest { version_id: original.version_id, delete_marker: false, - action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction, rule_id: "rule".to_string(), phase: crate::object_api::LifecycleDeleteAllPhase::Preflight, }), diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 6e8ba2e9b..2f3d87c7a 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -64,9 +64,9 @@ use futures::future::join_all; use http::HeaderMap; use lazy_static::lazy_static; use rand::RngExt as _; -use rustfs_common::heal_channel::{HealItemType, HealOpts}; use rustfs_config::server_config::Config; use rustfs_filemeta::FileInfo; +use rustfs_heal_contracts::heal_channel::{HealItemType, HealOpts}; use rustfs_lock::{LocalClient, LockClient, NamespaceLockWrapper}; use rustfs_madmin::heal_commands::HealResultItem; use rustfs_utils::path::{decode_dir_object, encode_dir_object, path_join_buf}; diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index 4521e087f..ddcd918ac 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -1049,7 +1049,7 @@ mod tests { lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest { version_id: Some(trigger_id), delete_marker: false, - action: rustfs_common::metrics::IlmAction::DeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction, rule_id: "rule".to_string(), phase: crate::object_api::LifecycleDeleteAllPhase::Preflight, }), @@ -1219,7 +1219,7 @@ mod tests { lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest { version_id: Some(marker_id), delete_marker: true, - action: rustfs_common::metrics::IlmAction::DelMarkerDeleteAllVersionsAction, + action: rustfs_scanner_contracts::metrics::IlmAction::DelMarkerDeleteAllVersionsAction, rule_id: "rule".to_string(), phase: crate::object_api::LifecycleDeleteAllPhase::Preflight, }), diff --git a/crates/ecstore/tests/ecstore_contract_compat_test.rs b/crates/ecstore/tests/ecstore_contract_compat_test.rs index d1d303f7a..ca20d42e5 100644 --- a/crates/ecstore/tests/ecstore_contract_compat_test.rs +++ b/crates/ecstore/tests/ecstore_contract_compat_test.rs @@ -14,8 +14,8 @@ mod storage_api; -use rustfs_common::heal_channel::HealOpts; use rustfs_filemeta::FileInfo; +use rustfs_heal_contracts::heal_channel::HealOpts; use rustfs_lock::NamespaceLockWrapper; use rustfs_madmin::heal_commands::HealResultItem; use storage_api::contract_compat::{ diff --git a/crates/heal/Cargo.toml b/crates/heal/Cargo.toml index ec2336fd7..ddedc83d4 100644 --- a/crates/heal/Cargo.toml +++ b/crates/heal/Cargo.toml @@ -77,6 +77,7 @@ rustfs-ecstore = { workspace = true } rustfs-lock = { workspace = true } rustfs-storage-api = { workspace = true } rustfs-common = { workspace = true } +rustfs-heal-contracts = { workspace = true } rustfs-madmin = { workspace = true } rustfs-utils = { workspace = true } tokio = { workspace = true, features = ["sync", "io-util", "time", "macros", "fs", "rt-multi-thread"] } diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index e00435100..7e2d38de0 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -19,7 +19,7 @@ use crate::heal::{ utils, }; use crate::{Error, Result}; -use rustfs_common::heal_channel::{ +use rustfs_heal_contracts::heal_channel::{ HealAdmissionReceipt, HealAdmissionResult, HealChannelCommand, HealChannelPriority, HealChannelReceiver, HealChannelRequest, HealChannelResponse, HealReceiptCommand, HealReceiptReceiver, HealRequestSource, HealScanMode, publish_heal_response, }; @@ -763,7 +763,7 @@ mod tests { use super::*; use crate::heal::manager::HealConfig; use crate::heal::storage::{HealObjectInfo, HealStorageAPI}; - use rustfs_common::heal_channel::{ + use rustfs_heal_contracts::heal_channel::{ HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest, HealRequestSource, HealScanMode, }; use std::sync::Arc; @@ -793,14 +793,14 @@ mod tests { _bucket: &str, _object: &str, _version_id: Option<&str>, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> crate::Result<(rustfs_madmin::heal_commands::HealResultItem, Option)> { Ok((rustfs_madmin::heal_commands::HealResultItem::default(), None)) } async fn heal_bucket( &self, _bucket: &str, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> crate::Result { Ok(rustfs_madmin::heal_commands::HealResultItem::default()) } @@ -1368,7 +1368,7 @@ mod tests { source: HealRequestSource::Admin, ..Default::default() }; - let mut responses = rustfs_common::heal_channel::subscribe_heal_responses(); + let mut responses = rustfs_heal_contracts::heal_channel::subscribe_heal_responses(); let (tx, rx) = oneshot::channel(); processor diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index 92a245c92..4d1e0b7c5 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -24,7 +24,7 @@ use crate::heal::{ use crate::{Error, Result}; use futures::{StreamExt, stream::FuturesUnordered}; use metrics::{counter, gauge}; -use rustfs_common::heal_channel::{HealOpts, HealRequestSource, HealScanMode}; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealRequestSource, HealScanMode}; use rustfs_madmin::heal_commands::HealResultItem; use std::sync::{ Arc, @@ -1346,7 +1346,7 @@ impl ErasureSetHealer { #[cfg(test)] mod tests { use super::{ErasureSetHealer, PageConcurrencyGuard}; - use rustfs_common::heal_channel::{HealRequestSource, HealScanMode}; + use rustfs_heal_contracts::heal_channel::{HealRequestSource, HealScanMode}; use std::sync::{ Arc, atomic::{AtomicUsize, Ordering}, @@ -1523,7 +1523,7 @@ mod resume_loop_tests { BUCKET_META_PREFIX, DiskOption, DiskStore, EcstoreError, Endpoint, HealDiskExt as _, RUSTFS_META_BUCKET, new_disk, }; use crate::{Error, Result}; - use rustfs_common::heal_channel::{HealOpts, HealRequestSource}; + use rustfs_heal_contracts::heal_channel::{HealOpts, HealRequestSource}; use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos}; use std::collections::{HashMap, HashSet, VecDeque}; use std::sync::atomic::{AtomicBool, Ordering}; diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index b6016c9e4..cb159772a 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -20,8 +20,10 @@ use crate::heal::{ }; use crate::{Error, Result}; use metrics::{counter, gauge}; -use rustfs_common::heal_channel::{HealAdmissionDropReason, HealAdmissionReceipt, HealAdmissionResult, HealRequestSource}; use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass}; +use rustfs_heal_contracts::heal_channel::{ + HealAdmissionDropReason, HealAdmissionReceipt, HealAdmissionResult, HealRequestSource, +}; use rustfs_madmin::heal_commands::HealResultItem; #[cfg(test)] use std::sync::LazyLock; diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 585b7a4e7..37a100f4b 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -17,8 +17,8 @@ use crate::heal::EcstoreError; use crate::heal::resume::{CheckpointManager, ReplacementTargetIdentity}; use crate::heal::storage::{HealObjectInfo, HealStorageAPI}; use crate::heal::task::{BatchHealFailure, HealOptions, HealPriority, HealRequest, HealTask, HealType}; -use rustfs_common::heal_channel::{HealOpts, HealRequestSource}; use rustfs_concurrency::{WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot}; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealRequestSource}; use rustfs_madmin::heal_commands::HealResultItem; use std::sync::Mutex as StdMutex; use tempfile::TempDir; diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 5cbe8e79b..026653d03 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -35,8 +35,8 @@ use super::{DiskStore, HealDiskExt as _, local_disk_map_read}; use crate::heal::manager::{HealManager, MrfRepairNoticeTarget}; use metrics::{counter, gauge}; -use rustfs_common::heal_channel::{HealAdmissionDropReason, HealAdmissionResult}; use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIntent}; +use rustfs_heal_contracts::heal_channel::{HealAdmissionDropReason, HealAdmissionResult}; use std::collections::{HashSet, VecDeque}; use std::sync::Arc; use std::time::Duration; @@ -483,7 +483,7 @@ pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest { options.set_index = usize::try_from(scope.set_index).ok(); } let mut request = HealRequest::new(heal_type, options, priority); - request.source = rustfs_common::heal_channel::HealRequestSource::Mrf; + request.source = rustfs_heal_contracts::heal_channel::HealRequestSource::Mrf; request } diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index 8c2c5fc2f..f84a75a60 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -16,7 +16,7 @@ use crate::{Error, Result}; use async_trait::async_trait; use base64::Engine as _; use base64::engine::general_purpose::URL_SAFE_NO_PAD; -use rustfs_common::heal_channel::{HealOpts, HealScanMode}; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; use rustfs_madmin::heal_commands::HealResultItem; use serde::{Deserialize, Serialize}; use std::sync::Arc; diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index c825869fe..ffd2ac618 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -23,8 +23,8 @@ use crate::heal::{ }; use crate::{Error, Result}; use metrics::{counter, histogram}; -use rustfs_common::heal_channel::{HealOpts, HealRequestSource, HealScanMode}; use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit}; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealRequestSource, HealScanMode}; use rustfs_madmin::heal_commands::HealResultItem; use rustfs_utils::path::SLASH_SEPARATOR; use serde::{Deserialize, Serialize}; diff --git a/crates/heal/src/lib.rs b/crates/heal/src/lib.rs index 7b91df47a..762c6c0c6 100644 --- a/crates/heal/src/lib.rs +++ b/crates/heal/src/lib.rs @@ -176,7 +176,7 @@ pub async fn init_heal_manager_with_workload_provider( let channel_receiver = if force_channel_failure { Err("forced heal channel initialization failure") } else { - rustfs_common::heal_channel::init_heal_channels() + rustfs_heal_contracts::heal_channel::init_heal_channels() }; let (receiver, receipt_receiver) = match channel_receiver { Ok(receivers) => receivers, @@ -358,7 +358,7 @@ mod tests { heal::storage::HealStorageAPI, init_heal_manager, run_owned_initialization, }; use crate::heal::storage_api::status::BucketInfo; - use rustfs_common::heal_channel::HealOpts; + use rustfs_heal_contracts::heal_channel::HealOpts; use rustfs_madmin::heal_commands::HealResultItem; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; diff --git a/crates/heal/tests/heal_b5_versioned_regression_test.rs b/crates/heal/tests/heal_b5_versioned_regression_test.rs index 05d789d4f..81cfdc79a 100644 --- a/crates/heal/tests/heal_b5_versioned_regression_test.rs +++ b/crates/heal/tests/heal_b5_versioned_regression_test.rs @@ -25,7 +25,6 @@ #![recursion_limit = "256"] use http::HeaderMap; -use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_heal::heal::{ manager::{HealConfig, HealManager}, storage::{ @@ -33,6 +32,7 @@ use rustfs_heal::heal::{ }, task::{HealOptions, HealPriority, HealRequest, HealTaskStatus, HealType}, }; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; use serial_test::serial; use std::{ path::{Path, PathBuf}, diff --git a/crates/heal/tests/heal_b920_subquorum_union_test.rs b/crates/heal/tests/heal_b920_subquorum_union_test.rs index eabda8e11..6be2d0595 100644 --- a/crates/heal/tests/heal_b920_subquorum_union_test.rs +++ b/crates/heal/tests/heal_b920_subquorum_union_test.rs @@ -24,10 +24,10 @@ #![recursion_limit = "256"] use http::HeaderMap; -use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_heal::heal::storage::{ ECStoreHealStorage, HealListItem, HealObjectOptions as ObjectOptions, HealPutObjReader as PutObjReader, HealStorageAPI, }; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; use serial_test::serial; use std::{ future::Future, diff --git a/crates/heal/tests/heal_bug_fixes_test.rs b/crates/heal/tests/heal_bug_fixes_test.rs index 59ff1bb3f..d2690b6be 100644 --- a/crates/heal/tests/heal_bug_fixes_test.rs +++ b/crates/heal/tests/heal_bug_fixes_test.rs @@ -121,14 +121,14 @@ fn test_heal_task_status_atomic_update() { _bucket: &str, _object: &str, _version_id: Option<&str>, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> rustfs_heal::Result<(rustfs_madmin::heal_commands::HealResultItem, Option)> { Ok((rustfs_madmin::heal_commands::HealResultItem::default(), None)) } async fn heal_bucket( &self, _bucket: &str, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> rustfs_heal::Result { Ok(rustfs_madmin::heal_commands::HealResultItem::default()) } @@ -221,7 +221,7 @@ async fn test_heal_task_transient_object_exists_skip_avoids_recreate() { _bucket: &str, _object: &str, _version_id: Option<&str>, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> rustfs_heal::Result<(rustfs_madmin::heal_commands::HealResultItem, Option)> { self.heal_object_calls.fetch_add(1, Ordering::SeqCst); Ok((rustfs_madmin::heal_commands::HealResultItem::default(), None)) @@ -230,7 +230,7 @@ async fn test_heal_task_transient_object_exists_skip_avoids_recreate() { async fn heal_bucket( &self, _bucket: &str, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> rustfs_heal::Result { Ok(rustfs_madmin::heal_commands::HealResultItem::default()) } diff --git a/crates/heal/tests/heal_integration_test.rs b/crates/heal/tests/heal_integration_test.rs index 940a69f73..213b43e86 100644 --- a/crates/heal/tests/heal_integration_test.rs +++ b/crates/heal/tests/heal_integration_test.rs @@ -15,12 +15,12 @@ #![recursion_limit = "256"] use http::HeaderMap; -use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_heal::heal::{ manager::{HealConfig, HealManager}, storage::{ECStoreHealStorage, HealObjectOptions as ObjectOptions, HealPutObjReader as PutObjReader, HealStorageAPI}, task::{HealOptions, HealPriority, HealRequest, HealTaskStatus, HealType}, }; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; use serial_test::serial; use std::{ path::{Path, PathBuf}, diff --git a/crates/lifecycle/Cargo.toml b/crates/lifecycle/Cargo.toml index 483b0e7d7..76f9d198c 100644 --- a/crates/lifecycle/Cargo.toml +++ b/crates/lifecycle/Cargo.toml @@ -57,6 +57,7 @@ hotpath.workspace = true async-trait.workspace = true metrics.workspace = true rustfs-common.workspace = true +rustfs-scanner-contracts.workspace = true rustfs-config = { workspace = true, features = ["constants"] } rustfs-replication.workspace = true rustfs-storage-api.workspace = true diff --git a/crates/lifecycle/src/core.rs b/crates/lifecycle/src/core.rs index 7dea47b71..2eadd6064 100644 --- a/crates/lifecycle/src/core.rs +++ b/crates/lifecycle/src/core.rs @@ -62,7 +62,7 @@ const ERR_LIFECYCLE_EXPIRED_OBJECT_DELETE_MARKER_WITH_TAGS: &str = const ERR_LIFECYCLE_RULE_MUST_HAVE_ACTION: &str = "Rule must have at least one of Expiration, Transition, NoncurrentVersionExpiration, NoncurrentVersionTransition, or DelMarkerExpiration"; const ERR_LIFECYCLE_PREFIX_FILTER_CONFLICT: &str = "Legacy Prefix and Filter cannot both be present in a lifecycle rule. Use Filter.Prefix instead of the top-level Prefix element."; -pub use rustfs_common::metrics::IlmAction; +pub use rustfs_scanner_contracts::metrics::IlmAction; #[async_trait::async_trait] pub trait RuleValidate { diff --git a/crates/lifecycle/src/evaluator.rs b/crates/lifecycle/src/evaluator.rs index 9ccc5b862..3b4bca449 100644 --- a/crates/lifecycle/src/evaluator.rs +++ b/crates/lifecycle/src/evaluator.rs @@ -18,8 +18,8 @@ use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration, ObjectLock use time::OffsetDateTime; use tracing::info; -use rustfs_common::metrics::IlmAction; use rustfs_replication::ReplicationStatusType; +use rustfs_scanner_contracts::metrics::IlmAction; use crate::object_lock; use crate::{Event, Lifecycle, ObjectOpts}; @@ -197,7 +197,7 @@ mod tests { use std::collections::HashMap; use std::sync::Arc; - use rustfs_common::metrics::IlmAction; + use rustfs_scanner_contracts::metrics::IlmAction; use s3s::dto::{ BucketLifecycleConfiguration, DefaultRetention, ExpirationStatus, LifecycleExpiration, LifecycleRule, NoncurrentVersionExpiration, ObjectLockConfiguration, ObjectLockEnabled, ObjectLockRetentionMode, ObjectLockRule, diff --git a/crates/lifecycle/src/lib.rs b/crates/lifecycle/src/lib.rs index 8986c9523..234dd9d7c 100644 --- a/crates/lifecycle/src/lib.rs +++ b/crates/lifecycle/src/lib.rs @@ -20,5 +20,5 @@ mod tagging; pub use core::*; pub use evaluator::Evaluator; -pub use rustfs_common::metrics::IlmAction; pub use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType}; +pub use rustfs_scanner_contracts::metrics::IlmAction; diff --git a/crates/obs/Cargo.toml b/crates/obs/Cargo.toml index 17ba113ed..338128814 100644 --- a/crates/obs/Cargo.toml +++ b/crates/obs/Cargo.toml @@ -105,6 +105,8 @@ workspace = true hotpath.workspace = true rustfs-audit = { workspace = true } rustfs-common = { workspace = true } +rustfs-heal-contracts = { workspace = true } +rustfs-scanner-contracts = { workspace = true } rustfs-config = { workspace = true, features = ["observability"] } # NOTE: This dependency on rustfs-ecstore is a known architectural limitation. # The obs crate imports types from ecstore for metrics collection. diff --git a/crates/obs/src/metrics/collectors/scanner.rs b/crates/obs/src/metrics/collectors/scanner.rs index 1e0221908..ed0b42d2c 100644 --- a/crates/obs/src/metrics/collectors/scanner.rs +++ b/crates/obs/src/metrics/collectors/scanner.rs @@ -493,7 +493,7 @@ mod tests { use super::*; use crate::metrics::report::report_metrics; use metrics_util::debugging::DebuggingRecorder; - use rustfs_common::metrics::{Metric, Metrics}; + use rustfs_scanner_contracts::metrics::{Metric, Metrics}; fn prometheus_counter_name(name: &str) -> String { if name.ends_with("_total") { diff --git a/crates/obs/src/metrics/stats_collector.rs b/crates/obs/src/metrics/stats_collector.rs index 0d3863762..8ee95ed47 100644 --- a/crates/obs/src/metrics/stats_collector.rs +++ b/crates/obs/src/metrics/stats_collector.rs @@ -37,16 +37,16 @@ use crate::metrics::{ }; use crate::node_identity::current_local_node_identity; use jiff::Timestamp; -use rustfs_common::heal_channel::HealScanMode; -use rustfs_common::metrics::{ - ScannerActiveBucketDriveSnapshot, ScannerBucketDriveResultSnapshot, ScannerMetricsReport, ScannerSourceWorkSnapshot, - global_metrics, -}; +use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_io_metrics::internode_metrics::global_internode_metrics; use rustfs_io_metrics::{ ProcessResourceSnapshot, ProcessSampler, ProcessStatusSnapshot, ProcessSystemSnapshot, s3_op_metrics_snapshot, snapshot_process_resource_and_system, snapshot_process_resource_and_system_with, }; +use rustfs_scanner_contracts::metrics::{ + ScannerActiveBucketDriveSnapshot, ScannerBucketDriveResultSnapshot, ScannerMetricsReport, ScannerSourceWorkSnapshot, + global_metrics, +}; use std::{ collections::{HashMap, HashSet}, sync::Arc, @@ -1667,7 +1667,7 @@ pub async fn collect_compression_cluster_stats() -> Option HealChannelRequest { diff --git a/crates/scanner/Cargo.toml b/crates/scanner/Cargo.toml index d941dce48..91f35715f 100644 --- a/crates/scanner/Cargo.toml +++ b/crates/scanner/Cargo.toml @@ -73,6 +73,8 @@ hotpath-cpu = [ hotpath.workspace = true rustfs-config = { workspace = true, features = ["server-config-model"] } rustfs-common = { workspace = true } +rustfs-heal-contracts = { workspace = true } +rustfs-scanner-contracts = { workspace = true } rustfs-credentials = { workspace = true } rustfs-utils = { workspace = true } tokio = { workspace = true, features = ["fs", "sync", "time", "macros", "rt-multi-thread"] } diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 084965989..8447c4917 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -23,7 +23,6 @@ use std::{ use http::HeaderMap; use metrics::{counter, describe_counter, describe_histogram, histogram}; -use rustfs_common::heal_channel::HealScanMode; #[cfg(test)] use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS; pub use rustfs_data_usage::{ @@ -33,6 +32,7 @@ pub use rustfs_data_usage::{ SizeReconciliationScope, SizeSummary, TierAccountingProof, TierStats, UNKNOWN_TIER, UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP, UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP, UnknownTierStats, hash_path, prefix_usage_in_cache, }; +use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use tokio::time::{Duration, Instant, sleep, timeout}; use tracing::{debug, warn}; diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 7f9f985da..2c78928b5 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -79,7 +79,7 @@ pub use remote_scanner::{ remote_scanner_request_matches_envelope, serve_remote_scanner_request, validate_remote_scanner_request_fence, }; pub use runtime_config::{apply_scanner_runtime_config, scanner_runtime_config_status, validate_scanner_runtime_config}; -pub use rustfs_common::last_minute; +pub use rustfs_scanner_contracts::last_minute; pub use scanner::{ ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerCycleScheduleStatus, init_data_scanner, reset_scanner_cycle_recovery, scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_topology_digest, diff --git a/crates/scanner/src/remote_scanner/stream.rs b/crates/scanner/src/remote_scanner/stream.rs index 594473196..7617f8fd0 100644 --- a/crates/scanner/src/remote_scanner/stream.rs +++ b/crates/scanner/src/remote_scanner/stream.rs @@ -26,9 +26,9 @@ use crate::{ scanner_publication_admission_for_epoch, scanner_publication_epoch, }; use hmac::{Hmac, KeyInit, Mac}; -use rustfs_common::heal_channel::HealScanMode; -use rustfs_common::metrics::{Metric, Metrics}; use rustfs_credentials::try_get_rpc_token; +use rustfs_heal_contracts::heal_channel::HealScanMode; +use rustfs_scanner_contracts::metrics::{Metric, Metrics}; use rustfs_utils::path::path_join_buf; use serde::{Deserialize, Serialize}; use sha2::Sha256; diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index e31cad163..d985931ad 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -40,12 +40,6 @@ use crate::{DataUsageInfo, ScannerActivityGuard, ScannerError, ScannerRuntimeGua use crate::{ScannerConfigObjectDelete, ScannerObjectIO, ScannerObjectOptions}; use bytes::Bytes; use chrono::{DateTime, Utc}; -use rustfs_common::heal_channel::HealScanMode; -use rustfs_common::metrics::{ - CurrentCycle, Metric, Metrics, ScanCyclePartialReason, ScanCycleWorkSnapshot, ScannerUsageSaveResult, ScannerWorkSource, - emit_scan_cycle_complete, emit_scan_cycle_deferred, emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded, - global_metrics, -}; use rustfs_config::ScannerSpeed; #[cfg(test)] use rustfs_config::{ @@ -54,7 +48,13 @@ use rustfs_config::{ }; use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS}; use rustfs_data_usage::observed_data_usage_is_newer; +use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_lock::{NamespaceLockGuard, error::LockError}; +use rustfs_scanner_contracts::metrics::{ + CurrentCycle, Metric, Metrics, ScanCyclePartialReason, ScanCycleWorkSnapshot, ScannerUsageSaveResult, ScannerWorkSource, + emit_scan_cycle_complete, emit_scan_cycle_deferred, emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded, + global_metrics, +}; use serde::{Deserialize, Serialize}; use sha2::{Digest as _, Sha256}; use tokio::sync::{Notify, mpsc}; diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index efc0ca26e..231bd147c 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -41,19 +41,19 @@ use crate::storage_api::owner::{ #[cfg(test)] use crate::storage_api::owner::{EcstoreExpirationStatus as ExpirationStatus, EcstoreLifecycleRule as LifecycleRule}; use metrics::{counter, describe_counter}; -use rustfs_common::heal_channel::{ - HEAL_DELETE_DANGLING, HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest, - HealRequestSource, HealScanMode, send_heal_request_with_admission, -}; -use rustfs_common::metrics::{ - CloseDiskGuard, IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, - UpdateCurrentPathFn, current_path_updater, global_metrics, -}; use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit, trace_subscriber_count}; use rustfs_filemeta::{ MAX_META_CACHE_HEAL_CANDIDATES, MAX_META_CACHE_HEAL_TRUNCATED_OBJECTS, MetaCacheEntries, MetaCacheEntry, MetaCacheHealCandidateKind, }; +use rustfs_heal_contracts::heal_channel::{ + HEAL_DELETE_DANGLING, HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest, + HealRequestSource, HealScanMode, send_heal_request_with_admission, +}; +use rustfs_scanner_contracts::metrics::{ + CloseDiskGuard, IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, + UpdateCurrentPathFn, current_path_updater, global_metrics, +}; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use time::OffsetDateTime; use tokio::select; diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 3da33b3b5..3708db8af 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -745,7 +745,7 @@ async fn test_scanner_heal_admission_accounting_maps_normal_scan_to_heal() { &metrics, HealScanMode::Normal, Ok(HealAdmissionResult::Dropped( - rustfs_common::heal_channel::HealAdmissionDropReason::QueueFull, + rustfs_heal_contracts::heal_channel::HealAdmissionDropReason::QueueFull, )), ); record_scanner_heal_admission(&metrics, HealScanMode::Normal, Err(())); @@ -1687,7 +1687,7 @@ fn test_describe_heal_admission_formats_unadmitted_results() { assert_eq!(describe_heal_admission(HealAdmissionResult::Full), "queue_full"); assert_eq!( describe_heal_admission(HealAdmissionResult::Dropped( - rustfs_common::heal_channel::HealAdmissionDropReason::QueueFull + rustfs_heal_contracts::heal_channel::HealAdmissionDropReason::QueueFull )), "dropped:queue_full" ); @@ -1879,10 +1879,10 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() { let healed_versions = Arc::new(Mutex::new(Vec::>::new())); let healed_versions_clone = healed_versions.clone(); let mut heal_rx = - rustfs_common::heal_channel::init_heal_channel().expect("heal channel should initialize once for scanner tests"); + rustfs_heal_contracts::heal_channel::init_heal_channel().expect("heal channel should initialize once for scanner tests"); let _heal_responder = tokio::spawn(async move { while let Some(command) = heal_rx.recv().await { - if let rustfs_common::heal_channel::HealChannelCommand::Start { + if let rustfs_heal_contracts::heal_channel::HealChannelCommand::Start { request, response_tx, .. } = command { diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 584dec769..8753b2139 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -24,13 +24,15 @@ use crate::{ use futures::future::join_all; use metrics::counter; use rand::seq::SliceRandom as _; -use rustfs_common::heal_channel::HealScanMode; -use rustfs_common::metrics::{Metric, Metrics, emit_scan_bucket_drive_complete, emit_scan_bucket_drive_partial, global_metrics}; #[cfg(test)] use rustfs_config::{ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS, ENV_SCANNER_MAX_CONCURRENT_SET_SCANS}; use rustfs_data_usage::{BucketTargetUsageInfo, BucketUsageInfo}; use rustfs_filemeta::FileMeta; +use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_lock::{LockError, NamespaceLockGuard}; +use rustfs_scanner_contracts::metrics::{ + Metric, Metrics, emit_scan_bucket_drive_complete, emit_scan_bucket_drive_partial, global_metrics, +}; use rustfs_utils::path::path_join_buf; use s3s::dto::{ BucketLifecycleConfiguration, ObjectLockConfiguration, ObjectLockEnabled, ReplicationConfiguration, VersioningConfiguration, diff --git a/crates/scanner/src/scanner_io/guards.rs b/crates/scanner/src/scanner_io/guards.rs index 3c93f759e..8e638a3c0 100644 --- a/crates/scanner/src/scanner_io/guards.rs +++ b/crates/scanner/src/scanner_io/guards.rs @@ -147,13 +147,13 @@ impl Drop for DiskBucketScanActiveGuard { pub(super) struct BucketDriveFailureGuard { failed: bool, - source: rustfs_common::metrics::ScannerWorkSource, + source: rustfs_scanner_contracts::metrics::ScannerWorkSource, bucket: String, drive: String, } impl BucketDriveFailureGuard { - pub(super) fn new(source: rustfs_common::metrics::ScannerWorkSource, bucket: &str, drive: &str) -> Self { + pub(super) fn new(source: rustfs_scanner_contracts::metrics::ScannerWorkSource, bucket: &str, drive: &str) -> Self { Self { failed: true, source, @@ -285,7 +285,7 @@ pub(super) fn scanner_task_join_error(stage: &str, err: tokio::task::JoinError) #[cfg(test)] mod tests { use super::*; - use rustfs_common::metrics::{ScannerWorkSource, global_metrics}; + use rustfs_scanner_contracts::metrics::{ScannerWorkSource, global_metrics}; #[test] fn bucket_drive_failure_guard_retires_active_scan_on_drop() { diff --git a/crates/scanner/src/scanner_io/io_disk.rs b/crates/scanner/src/scanner_io/io_disk.rs index cf209206e..1034d929e 100644 --- a/crates/scanner/src/scanner_io/io_disk.rs +++ b/crates/scanner/src/scanner_io/io_disk.rs @@ -161,8 +161,8 @@ impl ScannerIODisk for Disk { let bucket = cache.info.name.clone(); let disk_path = self.path().to_string_lossy().to_string(); let source = match scan_mode { - HealScanMode::Deep => rustfs_common::metrics::ScannerWorkSource::Bitrot, - HealScanMode::Normal | HealScanMode::Unknown => rustfs_common::metrics::ScannerWorkSource::Usage, + HealScanMode::Deep => rustfs_scanner_contracts::metrics::ScannerWorkSource::Bitrot, + HealScanMode::Normal | HealScanMode::Unknown => rustfs_scanner_contracts::metrics::ScannerWorkSource::Usage, }; global_metrics().record_scan_bucket_drive_start(source, &bucket, &disk_path); let mut failure_guard = BucketDriveFailureGuard::new(source, &bucket, &disk_path); diff --git a/crates/scanner/src/sleeper.rs b/crates/scanner/src/sleeper.rs index de46d53ef..1a554dcdd 100644 --- a/crates/scanner/src/sleeper.rs +++ b/crates/scanner/src/sleeper.rs @@ -16,11 +16,11 @@ use std::sync::atomic::{AtomicBool, AtomicU8, Ordering}; use std::sync::{Arc, LazyLock, RwLock}; use std::time::Instant; -use rustfs_common::metrics::global_metrics; use rustfs_config::{ DEFAULT_SCANNER_IDLE_MODE, DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS, ENV_SCANNER_IDLE_MODE, ENV_SCANNER_SPEED, ENV_SCANNER_YIELD_EVERY_N_OBJECTS, ScannerSpeed, }; +use rustfs_scanner_contracts::metrics::global_metrics; use tokio::time::Duration; const MIN_SLEEP: Duration = Duration::from_millis(1); diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index c392d5629..a37ebfef8 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -213,6 +213,8 @@ hotpath.workspace = true rustfs-heal = { workspace = true } rustfs-audit = { workspace = true } rustfs-common = { workspace = true } +rustfs-heal-contracts = { workspace = true } +rustfs-scanner-contracts = { workspace = true } rustfs-config = { workspace = true, features = ["notify", "server-config-model"] } rustfs-crypto = { workspace = true } rustfs-credentials = { workspace = true } diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index 28999fa94..57cdc2c6c 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -28,11 +28,11 @@ use futures_util::future::join_all; use http::{HeaderMap, HeaderValue, Uri}; use hyper::{Method, StatusCode}; use matchit::Params; -use rustfs_common::heal_channel::{ - HealAdmissionReceipt, HealChannelPriority, HealChannelRequest, HealOpts, HealRequestSource, HealScanMode, -}; use rustfs_config::MAX_HEAL_REQUEST_SIZE; use rustfs_heal::heal::utils::format_set_disk_id; +use rustfs_heal_contracts::heal_channel::{ + HealAdmissionReceipt, HealChannelPriority, HealChannelRequest, HealOpts, HealRequestSource, HealScanMode, +}; use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_scanner::scanner::{BackgroundHealInfo, read_background_heal_info}; use rustfs_utils::path::path_join; @@ -960,8 +960,8 @@ async fn submit_cluster_heal_start( } } -fn reject_heal_admission(result: rustfs_common::heal_channel::HealAdmissionResult) -> s3s::S3Error { - use rustfs_common::heal_channel::{HealAdmissionDropReason, HealAdmissionResult}; +fn reject_heal_admission(result: rustfs_heal_contracts::heal_channel::HealAdmissionResult) -> s3s::S3Error { + use rustfs_heal_contracts::heal_channel::{HealAdmissionDropReason, HealAdmissionResult}; match result { HealAdmissionResult::Full | HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull) => s3_error!( @@ -996,10 +996,10 @@ async fn submit_cluster_heal_channel_command( envelope: rustfs_protos::heal_control::Envelope, request_id: &str, response_id: String, -) -> S3Result { +) -> S3Result { match route_cluster_heal_control(&context, &route, envelope, request_id, false).await? { rustfs_protos::heal_control::Outcome::Channel { success, data, error } => { - Ok(rustfs_common::heal_channel::HealChannelResponse { + Ok(rustfs_heal_contracts::heal_channel::HealChannelResponse { request_id: response_id, success, data, @@ -1087,7 +1087,7 @@ fn build_heal_channel_request(hip: &HealInitParams) -> HealChannelRequest { } else { hip.hs.recursive }; - let mut heal_request = rustfs_common::heal_channel::create_heal_request( + let mut heal_request = rustfs_heal_contracts::heal_channel::create_heal_request( hip.bucket.clone(), if hip.obj_prefix.is_empty() { None @@ -1115,7 +1115,7 @@ fn build_heal_channel_request(hip: &HealInitParams) -> HealChannelRequest { } fn heal_channel_response_status( - response: &rustfs_common::heal_channel::HealChannelResponse, + response: &rustfs_heal_contracts::heal_channel::HealChannelResponse, ) -> (String, Vec, bool, Option) { let Some(data) = response.data.as_deref() else { return ("running".to_string(), Vec::new(), false, None); @@ -1136,19 +1136,21 @@ fn heal_channel_response_status( } #[cfg(test)] -fn heal_channel_response_summary(response: &rustfs_common::heal_channel::HealChannelResponse) -> String { +fn heal_channel_response_summary(response: &rustfs_heal_contracts::heal_channel::HealChannelResponse) -> String { heal_channel_response_status(response).0 } #[cfg(test)] fn heal_channel_response_items( - response: &rustfs_common::heal_channel::HealChannelResponse, + response: &rustfs_heal_contracts::heal_channel::HealChannelResponse, ) -> Vec { heal_channel_response_status(response).1 } #[cfg(test)] -fn heal_channel_response_progress(response: &rustfs_common::heal_channel::HealChannelResponse) -> Option { +fn heal_channel_response_progress( + response: &rustfs_heal_contracts::heal_channel::HealChannelResponse, +) -> Option { heal_channel_response_status(response).3 } @@ -1524,7 +1526,7 @@ mod tests { use http::StatusCode; use http::Uri; use matchit::Router; - use rustfs_common::heal_channel::{ + use rustfs_heal_contracts::heal_channel::{ HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealOpts, HealRequestSource, HealScanMode, }; use rustfs_scanner::scanner::BackgroundHealInfo; @@ -1718,7 +1720,7 @@ mod tests { dry_run: false, remove: true, recreate: false, - scan_mode: rustfs_common::heal_channel::HealScanMode::Normal, + scan_mode: rustfs_heal_contracts::heal_channel::HealScanMode::Normal, update_parity: false, no_lock: true, pool: Some(1), @@ -2640,11 +2642,15 @@ mod tests { #[test] fn test_heal_channel_response_summary_defaults_to_running() { - let response = rustfs_common::heal_channel::create_heal_response("token".to_string(), true, None, None); + let response = rustfs_heal_contracts::heal_channel::create_heal_response("token".to_string(), true, None, None); assert_eq!(heal_channel_response_summary(&response), "running"); - let response = - rustfs_common::heal_channel::create_heal_response("token".to_string(), true, Some(b"finished".to_vec()), None); + let response = rustfs_heal_contracts::heal_channel::create_heal_response( + "token".to_string(), + true, + Some(b"finished".to_vec()), + None, + ); assert_eq!(heal_channel_response_summary(&response), "finished"); } @@ -2668,7 +2674,7 @@ mod tests { "objectSize": 1024 }] }); - let response = rustfs_common::heal_channel::create_heal_response( + let response = rustfs_heal_contracts::heal_channel::create_heal_response( "token".to_string(), true, Some(serde_json::to_vec(&payload).expect("payload should serialize")), @@ -2695,7 +2701,7 @@ mod tests { "items": [], "progress": progress }); - let response = rustfs_common::heal_channel::create_heal_response( + let response = rustfs_heal_contracts::heal_channel::create_heal_response( "token".to_string(), true, Some(serde_json::to_vec(&payload).expect("payload should serialize")), diff --git a/rustfs/src/admin/handlers/scanner.rs b/rustfs/src/admin/handlers/scanner.rs index f932c3832..52fdfc777 100644 --- a/rustfs/src/admin/handlers/scanner.rs +++ b/rustfs/src/admin/handlers/scanner.rs @@ -24,10 +24,12 @@ use chrono::Utc; use http::{HeaderMap, HeaderValue}; use hyper::{Method, StatusCode}; use matchit::Params; -use rustfs_common::metrics::{ScannerLifecycleExpirySnapshot, ScannerMaintenanceControlSnapshot, ScannerMetricsReport}; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; use rustfs_credentials::Credentials; use rustfs_policy::policy::action::{Action, AdminAction}; +use rustfs_scanner_contracts::metrics::{ + ScannerLifecycleExpirySnapshot, ScannerMaintenanceControlSnapshot, ScannerMetricsReport, +}; use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index c4a597b08..3e2421f93 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -124,7 +124,7 @@ pub(crate) mod runtime_sources { pub(crate) type DailyAllTierStats = super::DailyAllTierStats; pub(crate) type ECStore = super::ECStore; pub(crate) type NotificationSys = super::NotificationSys; - pub(crate) type ScannerMetricsReport = rustfs_common::metrics::ScannerMetricsReport; + pub(crate) type ScannerMetricsReport = rustfs_scanner_contracts::metrics::ScannerMetricsReport; pub(crate) type StorageClassConfig = crate::storage::storage_api::ecstore_config::storageclass::Config; pub(crate) type TierConfigMgr = crate::storage::storage_api::TierConfigMgr; } diff --git a/rustfs/src/app/capacity_dirty_scope_test.rs b/rustfs/src/app/capacity_dirty_scope_test.rs index e00c21b73..3c84341a1 100644 --- a/rustfs/src/app/capacity_dirty_scope_test.rs +++ b/rustfs/src/app/capacity_dirty_scope_test.rs @@ -17,7 +17,7 @@ use super::storage_api::test::contract::bucket::{BucketOperations, BucketOptions use super::storage_api::test::contract::heal::HealOperations as _; use super::storage_api::test::contract::object::ObjectIO as _; use super::storage_api::test::{ECStore, Endpoint, EndpointServerPools, Endpoints, PoolEndpoints}; -use rustfs_common::heal_channel::{HealOpts, HealScanMode}; +use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; use rustfs_object_capacity::capacity_manager::{HybridStrategyConfig, create_isolated_manager}; use serial_test::serial; use std::{ diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index f11b10c38..c9b3438af 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -137,7 +137,7 @@ pub(crate) mod runtime { pub(crate) type NotificationSys = crate::storage::storage_api::NotificationSys; pub(crate) type ObjectStoreResolver = crate::storage::storage_api::ObjectStoreResolver; pub(crate) type ReplicationStats = crate::storage::storage_api::ReplicationStats; - pub(crate) type ScannerMetricsReport = rustfs_common::metrics::ScannerMetricsReport; + pub(crate) type ScannerMetricsReport = rustfs_scanner_contracts::metrics::ScannerMetricsReport; pub(crate) type StorageClassConfig = crate::storage::storage_api::ecstore_config::storageclass::Config; pub(crate) type TierConfigMgr = crate::storage::storage_api::TierConfigMgr; pub(crate) type TransitionState = crate::storage::storage_api::TransitionState; @@ -213,7 +213,7 @@ pub(crate) mod runtime { } pub(crate) async fn collect_scanner_metrics_report() -> ScannerMetricsReport { - rustfs_common::metrics::global_metrics().report().await + rustfs_scanner_contracts::metrics::global_metrics().report().await } #[cfg(test)] diff --git a/rustfs/src/cluster_snapshot.rs b/rustfs/src/cluster_snapshot.rs index e7535eb5d..f061fa981 100644 --- a/rustfs/src/cluster_snapshot.rs +++ b/rustfs/src/cluster_snapshot.rs @@ -23,9 +23,9 @@ use crate::storage_api::cluster::control_plane::{ ClusterPeerHealthSnapshot, ClusterPoolStateSnapshot, ClusterRpcBoundarySnapshot, }; use crate::workload_admission::workload_admission_registry_snapshot; -use rustfs_common::metrics::{ScannerMetricsReport, global_metrics}; use rustfs_concurrency::{AdmissionState, WorkloadAdmissionRegistrySnapshot}; use rustfs_io_metrics::internode_metrics::{InternodeMetricsSnapshot, global_internode_metrics}; +use rustfs_scanner_contracts::metrics::{ScannerMetricsReport, global_metrics}; #[derive(Debug, Clone, PartialEq, Eq)] pub struct ClusterReadOnlySnapshot { diff --git a/rustfs/src/startup_runtime_sources.rs b/rustfs/src/startup_runtime_sources.rs index 8b6e33115..ab63777b3 100644 --- a/rustfs/src/startup_runtime_sources.rs +++ b/rustfs/src/startup_runtime_sources.rs @@ -55,7 +55,7 @@ pub(crate) async fn publish_server_addr(addr: &str) { } pub(crate) async fn publish_init_time_now() { - rustfs_common::set_global_init_time_now().await; + rustfs_scanner_contracts::set_global_init_time_now().await; } pub(crate) fn init_kms_service_manager() -> Arc { diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index c73cf0f32..3f28c39ce 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -195,8 +195,8 @@ fn heal_control_remaining(expires_at_unix_ms: i64, now_unix_ms: i64) -> Result Result<(), Status> { - if request.source != rustfs_common::heal_channel::HealRequestSource::Admin { +fn validate_admin_heal_control_start(request: &rustfs_heal_contracts::heal_channel::HealChannelRequest) -> Result<(), Status> { + if request.source != rustfs_heal_contracts::heal_channel::HealRequestSource::Admin { return Err(Status::permission_denied("heal control start source must be admin")); } if request.pool_index.is_some() != request.set_index.is_some() { @@ -2542,7 +2542,7 @@ mod tests { _bucket: &str, _object: &str, _version_id: Option<&str>, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> rustfs_heal::Result<(rustfs_madmin::heal_commands::HealResultItem, Option)> { Ok((rustfs_madmin::heal_commands::HealResultItem::default(), None)) } @@ -2550,7 +2550,7 @@ mod tests { async fn heal_bucket( &self, _bucket: &str, - _opts: &rustfs_common::heal_channel::HealOpts, + _opts: &rustfs_heal_contracts::heal_channel::HealOpts, ) -> rustfs_heal::Result { Ok(rustfs_madmin::heal_commands::HealResultItem::default()) } @@ -2618,10 +2618,14 @@ mod tests { let now = i64::try_from(now).expect("test clock should fit in i64"); let metadata = || rustfs_protos::heal_control::RequestMetadata::new(rand::random(), now, now + 30_000, coordinator_epoch); let start = |request_id: String| { - let mut request = - rustfs_common::heal_channel::create_heal_request("bucket".to_string(), Some("prefix".to_string()), false, None); + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request( + "bucket".to_string(), + Some("prefix".to_string()), + false, + None, + ); request.id = request_id; - request.source = rustfs_common::heal_channel::HealRequestSource::Admin; + request.source = rustfs_heal_contracts::heal_channel::HealRequestSource::Admin; request }; @@ -3618,7 +3622,7 @@ mod tests { request } - let expired_request = rustfs_common::heal_channel::create_heal_request("bucket".to_string(), None, false, None); + let expired_request = rustfs_heal_contracts::heal_channel::create_heal_request("bucket".to_string(), None, false, None); let expired = rustfs_protos::heal_control::Envelope::start( expired_request, rustfs_protos::heal_control::RequestMetadata::new([1; 16], 1, 2, coordinator_epoch), @@ -3631,8 +3635,9 @@ mod tests { .expect_err("expired commands must fail before admission"); assert_eq!(expired.code(), tonic::Code::FailedPrecondition); - let mut non_admin_request = rustfs_common::heal_channel::create_heal_request("bucket".to_string(), None, false, None); - non_admin_request.source = rustfs_common::heal_channel::HealRequestSource::Scanner; + let mut non_admin_request = + rustfs_heal_contracts::heal_channel::create_heal_request("bucket".to_string(), None, false, None); + non_admin_request.source = rustfs_heal_contracts::heal_channel::HealRequestSource::Scanner; let now = OffsetDateTime::now_utc().unix_timestamp_nanos() / 1_000_000; let now = i64::try_from(now).expect("test clock should fit in i64"); let non_admin = rustfs_protos::heal_control::Envelope::start( diff --git a/rustfs/src/storage/rpc/node_service/heal.rs b/rustfs/src/storage/rpc/node_service/heal.rs index 8bbd27572..e5aaec78b 100644 --- a/rustfs/src/storage/rpc/node_service/heal.rs +++ b/rustfs/src/storage/rpc/node_service/heal.rs @@ -16,8 +16,8 @@ use crate::module_switches::{heal_enabled_from_env, scanner_enabled_from_env}; use crate::storage::storage_api::runtime_sources_consumer::EndpointServerPools; use jiff::Timestamp; use rmp_serde::Deserializer; -use rustfs_common::heal_channel::HealScanMode; use rustfs_heal::HealOperationsSnapshot; +use rustfs_heal_contracts::heal_channel::HealScanMode; use rustfs_scanner::scanner::BackgroundHealInfo; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; @@ -770,7 +770,7 @@ mod tests { BackgroundHealInfo { bitrot_start_time: Some(started_at), bitrot_start_cycle: 9, - current_scan_mode: rustfs_common::heal_channel::HealScanMode::Deep, + current_scan_mode: rustfs_heal_contracts::heal_channel::HealScanMode::Deep, }, HealOperationsSnapshot::default(), None, diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 5e85dc080..b9a8a5bfe 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -1545,7 +1545,7 @@ pub(crate) trait StoragePeerS3ClientExt { async fn heal_bucket_with_fence( &self, bucket: &str, - opts: &rustfs_common::heal_channel::HealOpts, + opts: &rustfs_heal_contracts::heal_channel::HealOpts, fenced_pools: &[usize], ) -> DiskResult; async fn make_bucket(&self, bucket: &str, opts: &contract::bucket::MakeBucketOptions) -> DiskResult<()>; @@ -1562,7 +1562,7 @@ impl StoragePeerS3ClientExt for LocalPeerS3Client { async fn heal_bucket_with_fence( &self, bucket: &str, - opts: &rustfs_common::heal_channel::HealOpts, + opts: &rustfs_heal_contracts::heal_channel::HealOpts, fenced_pools: &[usize], ) -> DiskResult { ecstore_rpc::PeerS3Client::heal_bucket_with_fence(self, bucket, opts, fenced_pools).await