refactor(replication): isolate status contract boundaries (#4236)

This commit is contained in:
Zhengchao An
2026-07-03 19:40:49 +08:00
committed by GitHub
parent 7001e53546
commit 769549dd2c
14 changed files with 119 additions and 49 deletions
@@ -22,6 +22,10 @@ use crate::bucket::lifecycle::lifecycle::{
self, Lifecycle, ObjectOpts, TransitionOptions, abort_incomplete_multipart_upload_due,
};
use crate::bucket::lifecycle::replication_sink;
use crate::bucket::lifecycle::replication_sink::{
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_statuses_map,
version_purge_statuses_map,
};
use crate::bucket::lifecycle::tier_delete_journal::{process_tier_delete_journal_entry, run_tier_delete_journal_recovery_loop};
use crate::bucket::lifecycle::tier_free_version_recovery::{DEFAULT_FREE_VERSION_RECOVERY_LIMIT, recover_tier_free_versions};
use crate::bucket::lifecycle::tier_last_day_stats::{DailyAllTierStats, LastDayTierStats};
@@ -57,10 +61,7 @@ use rustfs_config::{
ENV_TRANSITION_WORKERS_ABSOLUTE_MAX,
};
use rustfs_data_usage::TierStats;
use rustfs_filemeta::{
FileInfo, FileInfoOpts, NULL_VERSION_ID, ReplicateDecision, ReplicationState, RestoreStatusOps, VersionPurgeStatusType,
get_file_info, is_restored_object_on_disk,
};
use rustfs_filemeta::{FileInfo, FileInfoOpts, NULL_VERSION_ID, RestoreStatusOps, get_file_info, is_restored_object_on_disk};
use rustfs_utils::{get_env_i64, get_env_usize, path::encode_dir_object, string::strings_has_prefix_fold};
use s3s::dto::{
BucketLifecycleConfiguration, DefaultRetention, ExpirationStatus, ObjectLockConfiguration, RestoreRequest,
@@ -2890,12 +2891,12 @@ fn should_reuse_lifecycle_delete_replication_state(oi: &ObjectInfo, version_dele
if version_delete {
oi.version_purge_status == VersionPurgeStatusType::Pending && !state.purge_targets.is_empty()
} else {
oi.replication_status == rustfs_replication::ReplicationStatusType::Pending && !state.targets.is_empty()
oi.replication_status == ReplicationStatusType::Pending && !state.targets.is_empty()
}
}
fn lifecycle_version_purge_state_from_completed_targets(oi: &ObjectInfo) -> Option<ReplicationState> {
if oi.replication_status != rustfs_replication::ReplicationStatusType::Completed {
if oi.replication_status != ReplicationStatusType::Completed {
return None;
}
@@ -2909,7 +2910,7 @@ fn lifecycle_version_purge_state_from_completed_targets(oi: &ObjectInfo) -> Opti
Some(ReplicationState {
replicate_decision_str: oi.replication_decision.clone(),
version_purge_status_internal: Some(pending_status.clone()),
purge_targets: rustfs_replication::version_purge_statuses_map(&pending_status),
purge_targets: version_purge_statuses_map(&pending_status),
..Default::default()
})
}
@@ -2955,10 +2956,10 @@ fn replication_state_for_delete(dsc: ReplicateDecision, version_delete: bool) ->
};
if version_delete {
state.version_purge_status_internal = pending_status.clone();
state.purge_targets = rustfs_replication::version_purge_statuses_map(pending_status.as_deref().unwrap_or_default());
state.purge_targets = version_purge_statuses_map(pending_status.as_deref().unwrap_or_default());
} else {
state.replication_status_internal = pending_status.clone();
state.targets = rustfs_replication::replication_statuses_map(pending_status.as_deref().unwrap_or_default());
state.targets = replication_statuses_map(pending_status.as_deref().unwrap_or_default());
}
state
}
@@ -2998,6 +2999,9 @@ mod tests {
should_reuse_lifecycle_delete_replication_state, transitioned_cleanup_tuple,
};
use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc;
use crate::bucket::lifecycle::replication_sink::{
ReplicateDecision, ReplicateTargetDecision, ReplicationStatusType, VersionPurgeStatusType,
};
use crate::bucket::lifecycle::runtime_boundary as runtime_sources;
use crate::bucket::lifecycle::tier_sweeper::Jentry;
use crate::bucket::metadata::BUCKET_LIFECYCLE_CONFIG;
@@ -3017,7 +3021,6 @@ mod tests {
use futures::FutureExt;
use rustfs_common::metrics::{IlmAction, global_metrics};
use rustfs_config::ENV_TRANSITION_WORKERS_ABSOLUTE_MAX;
use rustfs_replication::{ReplicateDecision, ReplicationStatusType, VersionPurgeStatusType};
use s3s::dto::{
BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, MetadataEntry, OutputLocation,
RestoreRequest, RestoreRequestType, S3Location, Timestamp, Transition, TransitionStorageClass,
@@ -3915,7 +3918,7 @@ mod tests {
fn replication_state_for_delete_uses_replication_targets_for_current_delete() {
let arn = "arn:aws:s3:::target-bucket";
let mut dsc = ReplicateDecision::default();
dsc.set(rustfs_replication::ReplicateTargetDecision::new(arn.to_string(), true, false));
dsc.set(ReplicateTargetDecision::new(arn.to_string(), true, false));
let state = replication_state_for_delete(dsc, false);
@@ -3928,7 +3931,7 @@ mod tests {
fn replication_state_for_delete_uses_purge_targets_for_version_delete() {
let arn = "arn:aws:s3:::target-bucket";
let mut dsc = ReplicateDecision::default();
dsc.set(rustfs_replication::ReplicateTargetDecision::new(arn.to_string(), true, false));
dsc.set(ReplicateTargetDecision::new(arn.to_string(), true, false));
let state = replication_state_for_delete(dsc, true);
@@ -3953,7 +3956,7 @@ mod tests {
#[test]
fn lifecycle_delete_replication_state_does_not_reuse_put_replication_for_version_delete() {
let oi = ObjectInfo {
replication_status: rustfs_replication::ReplicationStatusType::Completed,
replication_status: ReplicationStatusType::Completed,
replication_status_internal: Some("arn:aws:s3:::target=COMPLETED;".to_string()),
replication_decision: "arn:aws:s3:::target=true;false;arn:aws:s3:::target;".to_string(),
..Default::default()
@@ -3968,7 +3971,7 @@ mod tests {
#[test]
fn lifecycle_version_purge_state_from_completed_targets_derives_pending_purge_targets() {
let oi = ObjectInfo {
replication_status: rustfs_replication::ReplicationStatusType::Completed,
replication_status: ReplicationStatusType::Completed,
replication_status_internal: Some("arn:aws:s3:::target=COMPLETED;".to_string()),
replication_decision: "arn:aws:s3:::target=true;false;arn:aws:s3:::target;".to_string(),
..Default::default()
+1 -1
View File
@@ -13,7 +13,6 @@
// limitations under the License.
use rustfs_config::{DEFAULT_ILM_PROCESS_TIME_SECS, ENV_ILM_PROCESS_TIME, ENV_ILM_PROCESS_TIME_DEPRECATED};
use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType};
use s3s::dto::{
BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, LifecycleRuleFilter,
NoncurrentVersionTransition, ObjectLockConfiguration, ObjectLockEnabled, RestoreRequest, Transition,
@@ -26,6 +25,7 @@ use time::{self, Duration, OffsetDateTime};
use tracing::debug;
use uuid::Uuid;
use crate::bucket::lifecycle::replication_sink::{ReplicationStatusType, VersionPurgeStatusType};
use crate::bucket::lifecycle::rule::{NoncurrentVersionTransitionOps, TransitionOps};
use crate::object_api::ObjectInfo;
@@ -170,7 +170,6 @@ mod tests {
use std::sync::Arc;
use rustfs_common::metrics::IlmAction;
use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType};
use s3s::dto::{
BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, Transition, TransitionStorageClass,
};
@@ -178,6 +177,7 @@ mod tests {
use uuid::Uuid;
use super::*;
use crate::bucket::lifecycle::replication_sink::{ReplicationStatusType, VersionPurgeStatusType};
fn expired_marker_lifecycle() -> Arc<BucketLifecycleConfiguration> {
Arc::new(BucketLifecycleConfiguration {
expiry_updated_at: None,
@@ -13,9 +13,14 @@
// limitations under the License.
use rustfs_common::metrics::IlmAction;
use rustfs_replication::{ReplicateDecision, ReplicationStatusType};
use crate::bucket::lifecycle::lifecycle::ObjectOpts;
#[cfg(test)]
pub(crate) use crate::bucket::replication::ReplicateTargetDecision;
pub(crate) use crate::bucket::replication::{
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_statuses_map,
version_purge_statuses_map,
};
use crate::bucket::replication::{ReplicationLifecycleBridge, ReplicationLifecycleConfig};
use crate::object_api::{ObjectInfo, ObjectOptions};
use crate::storage_api_contracts::object::{DeletedObject, ObjectToDelete};
@@ -70,7 +75,6 @@ mod tests {
use std::collections::HashMap;
use rustfs_common::metrics::IlmAction;
use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType};
use super::*;
@@ -44,6 +44,12 @@ mod runtime_boundary;
pub use config::{ObjectOpts, ReplicationConfigurationExt};
pub use datatypes::ResyncStatusType;
#[cfg(test)]
pub(crate) use replication_filemeta_boundary::ReplicateTargetDecision;
pub(crate) use replication_filemeta_boundary::{
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_statuses_map,
version_purge_statuses_map,
};
pub(crate) use replication_lifecycle_bridge::{ReplicationLifecycleBridge, ReplicationLifecycleConfig};
pub(crate) use replication_migration_bridge::ReplicationMigrationBridge;
pub use replication_object_bridge::ReplicationObjectBridge;
+1
View File
@@ -154,6 +154,7 @@ pub(crate) type ObjectOpts = EcstoreObjectOpts;
pub(crate) type ReplicationHealObject = ScannerReplicationHealObject;
pub(crate) type ReplicationHealQueueResult = ScannerReplicationHealResult;
pub(crate) type ReplicationQueueAdmission = ScannerReplicationQueueAdmission;
pub(crate) type ReplicationStatusType = storage_api::ReplicationStatusType;
pub(crate) type ScanGuard = EcstoreScanGuard;
pub(crate) type SetDisks = EcstoreSetDisks;
pub(crate) type StorageError = EcstoreStorageError;
+6 -6
View File
@@ -42,7 +42,6 @@ use rustfs_common::metrics::{
current_path_updater, global_metrics,
};
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;
@@ -53,10 +52,11 @@ use tracing::{debug, error, warn};
use crate::{
BucketVersioningSys, Disk, DiskError, DiskInfoOptions, Evaluator, Event, LcEventSrc, ListPathRawOptions, ObjectOpts,
ReplicationConfig, ReplicationHealObject, ReplicationQueueAdmission, ScannerDiskExt as _, ScannerLifecycleConfigExt as _,
ScannerVersioningConfigExt as _, StorageError, apply_expiry_rule, apply_transition_rule, enqueue_runtime_newer_noncurrent,
is_reserved_or_invalid_bucket, list_path_raw, path2_bucket_object, path2_bucket_object_with_base_path,
queue_replication_heal, scanner_is_erasure, scanner_replication_config_for_lifecycle_eval,
ReplicationConfig, ReplicationHealObject, ReplicationQueueAdmission, ReplicationStatusType, ScannerDiskExt as _,
ScannerLifecycleConfigExt as _, ScannerVersioningConfigExt as _, StorageError, apply_expiry_rule, apply_transition_rule,
enqueue_runtime_newer_noncurrent, is_reserved_or_invalid_bucket, list_path_raw, path2_bucket_object,
path2_bucket_object_with_base_path, queue_replication_heal, scanner_is_erasure,
scanner_replication_config_for_lifecycle_eval,
};
use crate::{ScannerObjectInfo as ObjectInfo, ScannerObjectToDelete as ObjectToDelete};
@@ -2862,9 +2862,9 @@ mod tests {
use crate::SCANNER_SLEEPER;
use super::*;
use crate::storage_api::VersionPurgeStatusType;
use crate::{DiskOption, Endpoint, new_disk};
use rustfs_filemeta::{FileInfo, FileMeta};
use rustfs_replication::VersionPurgeStatusType;
use serial_test::serial;
#[cfg(unix)]
use std::os::unix::fs::{PermissionsExt, symlink};
+2 -1
View File
@@ -15,7 +15,8 @@
use std::collections::HashMap;
use std::sync::Arc;
use rustfs_replication::{ReplicateObjectInfo, ReplicationStatusType, ReplicationType, VersionPurgeStatusType};
use rustfs_replication::{ReplicateObjectInfo, ReplicationType};
pub(crate) use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType};
use serde::{Deserialize, Serialize};
pub(crate) use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys as EcstoreBucketTargetSys;