mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-09 22:59:59 +00:00
refactor(replication): route app storage through ecstore (#4251)
This commit is contained in:
@@ -85,7 +85,6 @@ 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 }
|
||||
|
||||
@@ -609,18 +609,20 @@ pub(crate) mod bucket {
|
||||
use std::sync::Arc;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::storage::storage_api::ecstore_bucket::replication as replication_contracts;
|
||||
|
||||
pub(crate) type DeletedObjectReplicationInfo =
|
||||
crate::storage::storage_api::ecstore_bucket::replication::DeletedObjectReplicationInfo;
|
||||
type ReplicationObjectBridge = crate::storage::storage_api::ecstore_bucket::replication::ReplicationObjectBridge;
|
||||
pub(crate) type ReplicateDecision = rustfs_replication::ReplicateDecision;
|
||||
pub(crate) type ReplicateDecision = replication_contracts::ReplicateDecision;
|
||||
#[cfg(test)]
|
||||
pub(crate) type ReplicationState = rustfs_replication::ReplicationState;
|
||||
pub(crate) type ReplicationStatusType = rustfs_replication::ReplicationStatusType;
|
||||
pub(crate) type ReplicationTargetValidationError = rustfs_replication::ReplicationTargetValidationError;
|
||||
pub(crate) type VersionPurgeStatusType = rustfs_replication::VersionPurgeStatusType;
|
||||
pub(crate) const REPLICATE_INCOMING_DELETE: &str = rustfs_replication::REPLICATE_INCOMING_DELETE;
|
||||
pub(crate) type ReplicationState = replication_contracts::ReplicationState;
|
||||
pub(crate) type ReplicationStatusType = replication_contracts::ReplicationStatusType;
|
||||
pub(crate) type ReplicationTargetValidationError = replication_contracts::ReplicationTargetValidationError;
|
||||
pub(crate) type VersionPurgeStatusType = replication_contracts::VersionPurgeStatusType;
|
||||
pub(crate) const REPLICATE_INCOMING_DELETE: &str = replication_contracts::REPLICATE_INCOMING_DELETE;
|
||||
#[cfg(test)]
|
||||
pub(crate) use rustfs_replication::replication_statuses_map;
|
||||
pub(crate) use replication_contracts::replication_statuses_map;
|
||||
|
||||
pub(crate) async fn check_replicate_delete(
|
||||
bucket: &str,
|
||||
@@ -636,7 +638,7 @@ pub(crate) mod bucket {
|
||||
replication_source: &crate::storage::storage_api::StorageObjectInfo,
|
||||
deleted_delete_marker_version: bool,
|
||||
) -> Option<Uuid> {
|
||||
rustfs_replication::delete_replication_version_id(
|
||||
replication_contracts::delete_replication_version_id(
|
||||
replication_source.delete_marker,
|
||||
replication_source.version_id,
|
||||
deleted_delete_marker_version,
|
||||
@@ -655,7 +657,7 @@ pub(crate) mod bucket {
|
||||
user_defined,
|
||||
user_tags,
|
||||
status,
|
||||
rustfs_replication::ReplicationType::Object,
|
||||
replication_contracts::ReplicationType::Object,
|
||||
opts,
|
||||
);
|
||||
ReplicationObjectBridge::must_replicate(bucket, object, mopts).await
|
||||
@@ -666,7 +668,7 @@ pub(crate) mod bucket {
|
||||
store: Arc<crate::storage::storage_api::ECStore>,
|
||||
dsc: ReplicateDecision,
|
||||
) {
|
||||
ReplicationObjectBridge::schedule_object(oi, store, dsc, rustfs_replication::ReplicationType::Object).await;
|
||||
ReplicationObjectBridge::schedule_object(oi, store, dsc, replication_contracts::ReplicationType::Object).await;
|
||||
}
|
||||
|
||||
pub(crate) async fn schedule_replication_delete(dv: DeletedObjectReplicationInfo) {
|
||||
@@ -678,19 +680,19 @@ pub(crate) mod bucket {
|
||||
obj_info: &crate::storage::storage_api::StorageObjectInfo,
|
||||
version_id: Option<Uuid>,
|
||||
replica: bool,
|
||||
) -> Option<rustfs_replication::ReplicationState> {
|
||||
let source = rustfs_replication::ReplicationDeleteStateSource {
|
||||
) -> Option<replication_contracts::ReplicationState> {
|
||||
let source = replication_contracts::ReplicationDeleteStateSource {
|
||||
name: obj_info.name.clone(),
|
||||
user_tags: (*obj_info.user_tags).clone(),
|
||||
version_id,
|
||||
delete_marker: obj_info.delete_marker,
|
||||
replica,
|
||||
};
|
||||
rustfs_replication::delete_replication_state_from_config(config, &source)
|
||||
replication_contracts::delete_replication_state_from_config(config, &source)
|
||||
}
|
||||
|
||||
pub(crate) fn replication_target_arns(config: &s3s::dto::ReplicationConfiguration) -> HashSet<String> {
|
||||
rustfs_replication::replication_target_arns(config)
|
||||
replication_contracts::replication_target_arns(config)
|
||||
}
|
||||
|
||||
pub(crate) fn should_remove_replication_target(
|
||||
@@ -698,7 +700,7 @@ pub(crate) mod bucket {
|
||||
is_replication_service: bool,
|
||||
target_arns: &HashSet<String>,
|
||||
) -> bool {
|
||||
rustfs_replication::should_remove_replication_target(target_arn, is_replication_service, target_arns)
|
||||
replication_contracts::should_remove_replication_target(target_arn, is_replication_service, target_arns)
|
||||
}
|
||||
|
||||
pub(crate) fn should_schedule_delete_replication(
|
||||
@@ -706,7 +708,7 @@ pub(crate) mod bucket {
|
||||
replication_source: &crate::storage::storage_api::StorageObjectInfo,
|
||||
deleted_delete_marker_version: bool,
|
||||
) -> bool {
|
||||
rustfs_replication::should_schedule_delete_replication(rustfs_replication::ReplicationDeleteScheduleInput {
|
||||
replication_contracts::should_schedule_delete_replication(replication_contracts::ReplicationDeleteScheduleInput {
|
||||
replication_request: opts.replication_request,
|
||||
version_id_requested: opts.version_id.is_some(),
|
||||
source_delete_marker: replication_source.delete_marker,
|
||||
@@ -719,7 +721,7 @@ pub(crate) mod bucket {
|
||||
pub(crate) fn should_use_existing_delete_replication_info(
|
||||
opts: &crate::storage::storage_api::StorageObjectOptions,
|
||||
) -> bool {
|
||||
rustfs_replication::should_use_existing_delete_replication_info(opts.version_id.is_some(), opts.delete_marker)
|
||||
replication_contracts::should_use_existing_delete_replication_info(opts.version_id.is_some(), opts.delete_marker)
|
||||
}
|
||||
|
||||
pub(crate) fn should_use_existing_delete_replication_source(
|
||||
@@ -727,7 +729,7 @@ pub(crate) mod bucket {
|
||||
deleted_delete_marker: bool,
|
||||
has_existing_info: bool,
|
||||
) -> bool {
|
||||
rustfs_replication::should_use_existing_delete_replication_source(
|
||||
replication_contracts::should_use_existing_delete_replication_source(
|
||||
replication_request,
|
||||
deleted_delete_marker,
|
||||
has_existing_info,
|
||||
@@ -738,7 +740,7 @@ pub(crate) mod bucket {
|
||||
configured_arns: impl Iterator<Item = &'a str>,
|
||||
config: &s3s::dto::ReplicationConfiguration,
|
||||
) -> Result<(), ReplicationTargetValidationError> {
|
||||
rustfs_replication::validate_replication_config_target_arns(configured_arns, config)
|
||||
replication_contracts::validate_replication_config_target_arns(configured_arns, config)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user