From 18d89d1e318f35b7925313e0e3cf9e5217a86021 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Thu, 2 Jul 2026 08:48:42 +0800 Subject: [PATCH] refactor(replication): route metadata contracts through crate (#4160) --- Cargo.lock | 2 ++ crates/ecstore/src/bucket/bucket_target_sys.rs | 2 +- .../ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs | 2 +- crates/ecstore/src/bucket/lifecycle/core.rs | 2 +- crates/ecstore/src/bucket/lifecycle/evaluator.rs | 2 +- crates/ecstore/src/bucket/lifecycle/replication_sink.rs | 4 ++-- .../bucket/replication/replication_filemeta_boundary.rs | 4 ++-- crates/ecstore/src/client/object_handlers_common.rs | 2 +- crates/ecstore/src/object_api/types.rs | 3 ++- crates/ecstore/src/set_disk/mod.rs | 2 +- crates/replication/src/lib.rs | 7 +++++++ crates/scanner/Cargo.toml | 1 + crates/scanner/src/scanner_folder.rs | 6 ++++-- crates/scanner/src/storage_api.rs | 2 +- rustfs/Cargo.toml | 1 + rustfs/src/admin/router.rs | 2 +- rustfs/src/app/object_usecase.rs | 9 ++++----- rustfs/src/storage/options.rs | 4 ++-- 18 files changed, 35 insertions(+), 22 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7b5830252..20f1879cc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9122,6 +9122,7 @@ dependencies = [ "rustfs-policy", "rustfs-protocols", "rustfs-protos", + "rustfs-replication", "rustfs-rio", "rustfs-s3-ops", "rustfs-s3-types", @@ -10033,6 +10034,7 @@ dependencies = [ "rustfs-data-usage", "rustfs-ecstore", "rustfs-filemeta", + "rustfs-replication", "rustfs-storage-api", "rustfs-utils", "s3s", diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index d1a1a8b74..e940d21d1 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -49,7 +49,7 @@ use hyper_util::client::legacy::Client as HyperClient; use hyper_util::rt::{TokioExecutor, TokioTimer}; use reqwest::Client as HttpClient; use rustfs_config::{DEFAULT_TRUST_LEAF_CERT_AS_CA, ENV_TRUST_LEAF_CERT_AS_CA, RUSTFS_CA_CERT, RUSTFS_TLS_CERT}; -use rustfs_filemeta::ReplicationStatusType; +use rustfs_replication::ReplicationStatusType; use rustfs_utils::egress::{OutboundUrlError, validate_outbound_url}; use rustfs_utils::http::{ AMZ_BUCKET_REPLICATION_STATUS, AMZ_OBJECT_LOCK_BYPASS_GOVERNANCE, AMZ_OBJECT_LOCK_LEGAL_HOLD, AMZ_OBJECT_LOCK_MODE, diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 1906594b3..76306785a 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -3001,7 +3001,7 @@ mod tests { use futures::FutureExt; use rustfs_common::metrics::{IlmAction, global_metrics}; use rustfs_config::ENV_TRANSITION_WORKERS_ABSOLUTE_MAX; - use rustfs_filemeta::{ReplicateDecision, ReplicationStatusType, VersionPurgeStatusType}; + use rustfs_replication::{ReplicateDecision, ReplicationStatusType, VersionPurgeStatusType}; use s3s::dto::{ BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, MetadataEntry, OutputLocation, RestoreRequest, RestoreRequestType, S3Location, Timestamp, Transition, TransitionStorageClass, diff --git a/crates/ecstore/src/bucket/lifecycle/core.rs b/crates/ecstore/src/bucket/lifecycle/core.rs index bfcf921bb..c662e0b25 100644 --- a/crates/ecstore/src/bucket/lifecycle/core.rs +++ b/crates/ecstore/src/bucket/lifecycle/core.rs @@ -13,7 +13,7 @@ // limitations under the License. use rustfs_config::{DEFAULT_ILM_PROCESS_TIME_SECS, ENV_ILM_PROCESS_TIME, ENV_ILM_PROCESS_TIME_DEPRECATED}; -use rustfs_filemeta::{ReplicationStatusType, VersionPurgeStatusType}; +use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType}; use s3s::dto::{ BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, LifecycleRuleFilter, NoncurrentVersionTransition, ObjectLockConfiguration, ObjectLockEnabled, RestoreRequest, Transition, diff --git a/crates/ecstore/src/bucket/lifecycle/evaluator.rs b/crates/ecstore/src/bucket/lifecycle/evaluator.rs index 30f8d2315..fdb199438 100644 --- a/crates/ecstore/src/bucket/lifecycle/evaluator.rs +++ b/crates/ecstore/src/bucket/lifecycle/evaluator.rs @@ -170,7 +170,7 @@ mod tests { use std::sync::Arc; use rustfs_common::metrics::IlmAction; - use rustfs_filemeta::{ReplicationStatusType, VersionPurgeStatusType}; + use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType}; use s3s::dto::{ BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, Transition, TransitionStorageClass, }; diff --git a/crates/ecstore/src/bucket/lifecycle/replication_sink.rs b/crates/ecstore/src/bucket/lifecycle/replication_sink.rs index 24c1d8ab7..85a06032a 100644 --- a/crates/ecstore/src/bucket/lifecycle/replication_sink.rs +++ b/crates/ecstore/src/bucket/lifecycle/replication_sink.rs @@ -13,7 +13,7 @@ // limitations under the License. use rustfs_common::metrics::IlmAction; -use rustfs_filemeta::{ReplicateDecision, ReplicationStatusType}; +use rustfs_replication::{ReplicateDecision, ReplicationStatusType}; use crate::bucket::lifecycle::lifecycle::ObjectOpts; use crate::bucket::replication::{ReplicationLifecycleBridge, ReplicationLifecycleConfig}; @@ -70,7 +70,7 @@ mod tests { use std::collections::HashMap; use rustfs_common::metrics::IlmAction; - use rustfs_filemeta::{ReplicationStatusType, VersionPurgeStatusType}; + use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType}; use super::*; diff --git a/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs b/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs index 80622d2a6..01d8a5ce9 100644 --- a/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs @@ -12,11 +12,11 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub(crate) use rustfs_filemeta::{ +pub(crate) use rustfs_replication::{MrfOpKind, MrfReplicateEntry}; +pub(crate) use rustfs_replication::{ REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicateObjectInfo, ReplicateTargetDecision, ReplicatedInfos, ReplicatedTargetInfo, ReplicationAction, ReplicationState, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation, ResyncDecision, ResyncTargetDecision, VersionPurgeStatusType, get_replication_state, parse_replicate_decision, replication_statuses_map, target_reset_header, version_purge_statuses_map, }; -pub(crate) use rustfs_replication::{MrfOpKind, MrfReplicateEntry}; diff --git a/crates/ecstore/src/client/object_handlers_common.rs b/crates/ecstore/src/client/object_handlers_common.rs index 7a8e61e7e..b08a7683f 100644 --- a/crates/ecstore/src/client/object_handlers_common.rs +++ b/crates/ecstore/src/client/object_handlers_common.rs @@ -27,8 +27,8 @@ use crate::bucket::versioning_sys::BucketVersioningSys; use crate::object_api::ObjectOptions; use crate::storage_api_contracts::object::{ObjectOperations as _, ObjectToDelete}; use crate::store::ECStore; -use rustfs_filemeta::ReplicationState; use rustfs_lock::MAX_DELETE_LIST; +use rustfs_replication::ReplicationState; pub async fn delete_object_versions(api: &Arc, bucket: &str, to_del: &[ObjectToDelete], _lc_event: lifecycle::Event) { let version_suspended = match BucketVersioningSys::get(bucket).await { diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 57bde6189..6ef5724d1 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -790,7 +790,8 @@ fn versions_after_marker(file_infos: &rustfs_filemeta::FileInfoVersions, marker: #[cfg(test)] mod tests { use super::*; - use rustfs_filemeta::{FileInfo, FileMeta, MetaCacheEntry, ReplicationState, TRANSITION_COMPLETE}; + use rustfs_filemeta::{FileInfo, FileMeta, MetaCacheEntry, TRANSITION_COMPLETE}; + use rustfs_replication::ReplicationState; #[test] fn versions_after_marker_handles_null_version_marker() { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 35165e133..5811e8e40 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -6847,9 +6847,9 @@ mod tests { use rustfs_filemeta::ErasureInfo; use rustfs_filemeta::FileMeta; use rustfs_filemeta::MetaCacheEntry; - use rustfs_filemeta::ReplicationState; use rustfs_lock::client::local::LocalClient; use rustfs_lock::{LockError, LockInfo, LockResponse, LockStats}; + use rustfs_replication::ReplicationState; use serial_test::serial; use std::collections::HashMap; use tempfile::TempDir; diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index 0fc000147..2d809bc10 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -20,3 +20,10 @@ pub use resync::{ BucketReplicationResyncStatus, Error, Result, ResyncOpts, ResyncStatusType, TargetReplicationResyncStatus, decode_resync_file, encode_resync_file, }; +pub use rustfs_filemeta::{ + REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING_DELETE, + ReplicateDecision, ReplicateObjectInfo, ReplicateTargetDecision, ReplicatedInfos, ReplicatedTargetInfo, ReplicationAction, + ReplicationState, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation, ResyncDecision, ResyncTargetDecision, + VersionPurgeStatusType, get_replication_state, parse_replicate_decision, replication_statuses_map, target_reset_header, + version_purge_statuses_map, +}; diff --git a/crates/scanner/Cargo.toml b/crates/scanner/Cargo.toml index 95389a997..d50fe992b 100644 --- a/crates/scanner/Cargo.toml +++ b/crates/scanner/Cargo.toml @@ -44,6 +44,7 @@ time = { workspace = true } chrono = { workspace = true } rmp-serde = { workspace = true } rustfs-filemeta = { workspace = true } +rustfs-replication = { workspace = true } tokio-util = { workspace = true } rustfs-ecstore = { workspace = true } rustfs-storage-api = { workspace = true } diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 8f8366b7a..5f8930255 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -41,7 +41,8 @@ use rustfs_common::metrics::{ IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, UpdateCurrentPathFn, current_path_updater, global_metrics, }; -use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ReplicationStatusType}; +use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams}; +use rustfs_replication::ReplicationStatusType; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration}; use time::OffsetDateTime; @@ -2836,7 +2837,8 @@ mod tests { use super::*; use crate::{DiskOption, Endpoint, new_disk}; - use rustfs_filemeta::{FileInfo, FileMeta, VersionPurgeStatusType}; + use rustfs_filemeta::{FileInfo, FileMeta}; + use rustfs_replication::VersionPurgeStatusType; use serial_test::serial; #[cfg(unix)] use std::os::unix::fs::{PermissionsExt, symlink}; diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index a31e865a1..545f823d3 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -15,7 +15,7 @@ use std::collections::HashMap; use std::sync::Arc; -use rustfs_filemeta::{ReplicateObjectInfo, ReplicationStatusType, ReplicationType, VersionPurgeStatusType}; +use rustfs_replication::{ReplicateObjectInfo, ReplicationStatusType, ReplicationType, VersionPurgeStatusType}; use serde::{Deserialize, Serialize}; pub(crate) use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys as EcstoreBucketTargetSys; diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 9a3da8c86..e5f29b4e1 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -85,6 +85,7 @@ rustfs-policy = { workspace = true } rustfs-protocols = { workspace = true } rustfs-protos = { workspace = true } rustfs-rio = { workspace = true } +rustfs-replication = { workspace = true } rustfs-s3-types = { workspace = true } rustfs-s3-ops = { workspace = true } rustfs-security-governance = { workspace = true } diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index c1563dac5..14d264a82 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -59,10 +59,10 @@ use rustfs_config::{ ENABLE_KEY, WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CA, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY, WEBHOOK_ENDPOINT, WEBHOOK_SKIP_TLS_VERIFY, }; -use rustfs_filemeta::ReplicationStatusType; use rustfs_madmin::utils::parse_duration; use rustfs_notify::{Event as NotificationEvent, notification_system}; use rustfs_policy::policy::action::{Action, S3Action}; +use rustfs_replication::ReplicationStatusType; use rustfs_s3_types::EventName; use rustfs_signer::pre_sign_v4; use rustfs_utils::egress::validate_outbound_url; diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index cb74abad7..9de95333f 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -109,17 +109,16 @@ use metrics::{counter, histogram}; use pin_project_lite::pin_project; use rustfs_concurrency::GetObjectQueueSnapshot; use rustfs_config::MI_B; -use rustfs_filemeta::{ - REPLICATE_INCOMING_DELETE, ReplicationStatusType, RestoreStatusOps, VersionPurgeStatusType, parse_restore_obj_status, -}; -#[cfg(test)] -use rustfs_filemeta::{ReplicationState, replication_statuses_map}; +use rustfs_filemeta::{RestoreStatusOps, parse_restore_obj_status}; use rustfs_io_core::{BytesPool, PooledBuffer}; use rustfs_io_metrics; use rustfs_lock::NamespaceLockGuard; use rustfs_notify::EventArgsBuilder; use rustfs_object_capacity::capacity_manager::get_capacity_manager; use rustfs_policy::policy::action::{Action, S3Action}; +use rustfs_replication::{REPLICATE_INCOMING_DELETE, ReplicationStatusType, VersionPurgeStatusType}; +#[cfg(test)] +use rustfs_replication::{ReplicationState, replication_statuses_map}; use rustfs_s3_ops::{S3Operation, delete_event_name_for_marker, put_event_name_for_post_object}; use rustfs_s3select_api::object_store::bytes_stream; use rustfs_targets::{ diff --git a/rustfs/src/storage/options.rs b/rustfs/src/storage/options.rs index 278539d58..f13df7d7d 100644 --- a/rustfs/src/storage/options.rs +++ b/rustfs/src/storage/options.rs @@ -16,7 +16,7 @@ use super::{BucketVersioningSys, Result, StorageError}; use crate::storage::storage_api::options_consumer::contract::{object::HTTPPreconditions, range::HTTPRangeSpec}; use http::header::{IF_MATCH, IF_NONE_MATCH}; use http::{HeaderMap, HeaderValue}; -use rustfs_filemeta::ReplicationStatusType; +use rustfs_replication::ReplicationStatusType; use rustfs_utils::http::{ AMZ_BUCKET_REPLICATION_STATUS, SUFFIX_FORCE_DELETE, SUFFIX_REPLICATION_ACTUAL_OBJECT_SIZE, SUFFIX_REPLICATION_SSEC_CRC, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_REQUEST, SUFFIX_SOURCE_VERSION_ID, get_header, @@ -820,7 +820,7 @@ mod tests { get_default_opts, get_opts, parse_copy_source_range, put_opts, put_opts_from_headers, validate_archive_content_encoding, }; use http::{HeaderMap, HeaderValue}; - use rustfs_filemeta::ReplicationStatusType; + use rustfs_replication::ReplicationStatusType; use rustfs_utils::http::{ AMZ_BUCKET_REPLICATION_STATUS, AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER, AMZ_OBJECT_LOCK_MODE_LOWER, AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER, SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_REQUEST,