refactor(replication): own filemeta wire contracts (#4254)

This commit is contained in:
Zhengchao An
2026-07-04 06:39:54 +08:00
committed by GitHub
parent 09fcacbe2b
commit 3d4fb532e5
22 changed files with 1196 additions and 86 deletions
@@ -23,8 +23,8 @@ use crate::bucket::lifecycle::lifecycle::{
};
use crate::bucket::lifecycle::replication_sink;
use crate::bucket::lifecycle::replication_sink::{
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_statuses_map,
version_purge_statuses_map,
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta,
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};
@@ -2881,7 +2881,7 @@ async fn schedule_lifecycle_replication_delete_if_needed(oi: &ObjectInfo, dobj:
return;
}
delete_object.replication_state = replication_state;
delete_object.replication_state = replication_state.as_ref().map(replication_state_to_filemeta);
replication_sink::schedule_delete(oi.bucket.clone(), delete_object).await;
}
@@ -18,8 +18,8 @@ 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,
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta,
replication_statuses_map, version_purge_statuses_map,
};
use crate::bucket::replication::{ReplicationLifecycleBridge, ReplicationLifecycleConfig};
use crate::object_api::{ObjectInfo, ObjectOptions};
+10 -11
View File
@@ -37,11 +37,11 @@ paths.
| `EcstoreReplicationBoundaryImports` | ECStore-side imports from `rustfs-replication`. | Direct `rustfs-replication` imports under `crates/ecstore/src/bucket/replication` stay in `*_boundary.rs` modules, including config and resync facade re-exports. |
| `RuntimeReplicationFacadeConsumers` | Runtime owner consumers of replication DTOs and status types. | Scanner, admin, and storage owner facades import replication DTOs/status types through `rustfs-ecstore`; app storage keeps the remaining direct object/delete helper calls behind its local storage API boundary. |
| `ReplicationResyncContracts` | Resync options, target status, bucket status, status classifiers, and persisted resync/MRF status wire format. | Owned by `crates/replication`; ECStore imports them through `replication_resync_boundary.rs`, which maps crate errors to ECStore errors. |
| `ReplicationCrateFileMetaFacade` | Replication facade compatibility symbols that still originate in filemeta wire contracts. | `crates/replication/src/filemeta.rs` is the only direct `rustfs-filemeta` import boundary inside `rustfs-replication`. |
| `ReplicationCrateFileMetaIndependence` | Replication status, decision, MRF, resync, and target-reset wire contracts owned by `rustfs-replication`. | `crates/replication/src/filemeta.rs` owns these contracts; `rustfs-replication` must not import or depend on `rustfs-filemeta`. |
| `ReplicationConfigStore` | Replication config persistence and config-derived labels used by target options. | Config read/save helpers and storage class labels are exposed through the contract type in `replication_config_store.rs`. |
| `ReplicationFileMeta` | Replication status, decisions, MRF entries, resync decisions, and target reset helpers. | `rustfs_filemeta` replication contracts are concentrated in `replication_filemeta_boundary.rs`; `FileInfo` remains in the storage boundary for storage trait bindings and walk options. |
| `StorageApiReplicationContracts` | Storage-api delete DTO replication state/status helpers. | The temporary `rustfs-filemeta` dependency is isolated in `crates/storage-api/src/replication.rs` until the wire contracts can move without creating a `rustfs-replication` / `rustfs-storage-api` cycle. |
| `ReplicationCrateStorageApiBoundary` | Storage API delete DTOs consumed by `rustfs-replication`. | `crates/replication/src/storage_api.rs` is the only direct `rustfs-storage-api` import boundary inside `rustfs-replication`. |
| `ReplicationFileMeta` | ECStore compatibility conversions for filemeta replication state/status. | `rustfs_filemeta` to `rustfs_replication` conversions are concentrated in `replication_filemeta_boundary.rs`; `FileInfo` remains in the storage boundary for storage trait bindings and walk options. |
| `StorageApiReplicationContracts` | Storage-api delete DTO replication state/status helpers. | Storage-api owner DTOs keep their local replication boundary; ECStore converts them in `replication_storage_boundary.rs` before queueing replication work. |
| `ReplicationCrateStorageApiIndependence` | Delete work DTOs consumed by `rustfs-replication`. | `crates/replication/src/storage_api.rs` owns these DTOs; `rustfs-replication` must not import or depend on `rustfs-storage-api`. |
| `ReplicationObjectDecisionContracts` | Object replication options, delete replication decisions, resync target projection, multipart planning, and delete-marker retry classifiers. | Owned by `crates/replication`; ECStore imports them through `replication_object_decision_boundary.rs`. |
| `ReplicationQueueContracts` | Queue admission, heal queue results/actions, worker operations, worker sizing, and backpressure decisions. | Owned by `crates/replication`; ECStore imports them through `replication_queue_boundary.rs`. |
| `ReplicationStatsContracts` | Bucket stats, replication target stats, queue/proxy metrics, and worker metric snapshots. | Owned by `crates/replication`; ECStore imports them through `replication_stats_boundary.rs`. |
@@ -88,13 +88,12 @@ paths.
10. Keep ECStore owner modules outside `bucket/replication` behind bridge
contracts when they need replication codec or config helper behavior.
11. Keep storage-api replication status/state helpers behind
`crates/storage-api/src/replication.rs` until the underlying wire contracts
can move without a `rustfs-replication` / `rustfs-storage-api` dependency
cycle.
12. Keep direct `rustfs-filemeta` imports inside `rustfs-replication`
concentrated in `crates/replication/src/filemeta.rs`.
13. Keep direct `rustfs-storage-api` imports inside `rustfs-replication`
concentrated in `crates/replication/src/storage_api.rs`.
`crates/storage-api/src/replication.rs`; ECStore converts owner DTOs at the
replication storage boundary.
12. Keep `rustfs-replication` independent from `rustfs-filemeta`; ECStore
compatibility conversions live in `replication_filemeta_boundary.rs`.
13. Keep `rustfs-replication` independent from `rustfs-storage-api`; ECStore
compatibility conversions live in `replication_storage_boundary.rs`.
14. Keep direct `rustfs-replication` imports inside ECStore replication
concentrated in `*_boundary.rs` modules.
15. Keep scanner, admin, and storage-owner replication status/DTO consumers
+5 -1
View File
@@ -52,7 +52,11 @@ pub(crate) use replication_filemeta_boundary::ReplicateTargetDecision;
pub(crate) use replication_filemeta_boundary::version_purge_statuses_map;
pub use replication_filemeta_boundary::{
REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicateObjectInfo, ReplicationState, ReplicationStatusType, ReplicationType,
VersionPurgeStatusType, replication_statuses_map,
VersionPurgeStatusType, replication_state_to_filemeta, replication_status_to_filemeta, replication_statuses_map,
version_purge_status_to_filemeta,
};
pub(crate) use replication_filemeta_boundary::{
replication_state_from_filemeta, replication_status_from_filemeta, version_purge_status_from_filemeta,
};
pub(crate) use replication_lifecycle_bridge::{ReplicationLifecycleBridge, ReplicationLifecycleConfig};
pub(crate) use replication_migration_bridge::ReplicationMigrationBridge;
@@ -22,3 +22,65 @@ pub use rustfs_replication::{
REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicateObjectInfo, ReplicationState, ReplicationStatusType, ReplicationType,
VersionPurgeStatusType, replication_statuses_map,
};
pub(crate) fn replication_status_from_filemeta(status: rustfs_filemeta::ReplicationStatusType) -> ReplicationStatusType {
ReplicationStatusType::from(status.as_str())
}
pub(crate) fn version_purge_status_from_filemeta(status: rustfs_filemeta::VersionPurgeStatusType) -> VersionPurgeStatusType {
VersionPurgeStatusType::from(status.as_str())
}
pub(crate) fn replication_state_from_filemeta(state: &rustfs_filemeta::ReplicationState) -> ReplicationState {
ReplicationState {
replica_timestamp: state.replica_timestamp,
replica_status: replication_status_from_filemeta(state.replica_status.clone()),
delete_marker: state.delete_marker,
replication_timestamp: state.replication_timestamp,
replication_status_internal: state.replication_status_internal.clone(),
version_purge_status_internal: state.version_purge_status_internal.clone(),
replicate_decision_str: state.replicate_decision_str.clone(),
targets: state
.targets
.iter()
.map(|(arn, status)| (arn.clone(), replication_status_from_filemeta(status.clone())))
.collect(),
purge_targets: state
.purge_targets
.iter()
.map(|(arn, status)| (arn.clone(), version_purge_status_from_filemeta(status.clone())))
.collect(),
reset_statuses_map: state.reset_statuses_map.clone(),
}
}
pub fn replication_status_to_filemeta(status: ReplicationStatusType) -> rustfs_filemeta::ReplicationStatusType {
rustfs_filemeta::ReplicationStatusType::from(status.as_str())
}
pub fn version_purge_status_to_filemeta(status: VersionPurgeStatusType) -> rustfs_filemeta::VersionPurgeStatusType {
rustfs_filemeta::VersionPurgeStatusType::from(status.as_str())
}
pub fn replication_state_to_filemeta(state: &ReplicationState) -> rustfs_filemeta::ReplicationState {
rustfs_filemeta::ReplicationState {
replica_timestamp: state.replica_timestamp,
replica_status: replication_status_to_filemeta(state.replica_status.clone()),
delete_marker: state.delete_marker,
replication_timestamp: state.replication_timestamp,
replication_status_internal: state.replication_status_internal.clone(),
version_purge_status_internal: state.version_purge_status_internal.clone(),
replicate_decision_str: state.replicate_decision_str.clone(),
targets: state
.targets
.iter()
.map(|(arn, status)| (arn.clone(), replication_status_to_filemeta(status.clone())))
.collect(),
purge_targets: state
.purge_targets
.iter()
.map(|(arn, status)| (arn.clone(), version_purge_status_to_filemeta(status.clone())))
.collect(),
reset_statuses_map: state.reset_statuses_map.clone(),
}
}
@@ -16,6 +16,7 @@ use rustfs_filemeta::FileInfo;
use tokio_util::sync::CancellationToken;
use super::replication_error_boundary::Error;
use super::replication_filemeta_boundary::{replication_state_from_filemeta, version_purge_status_from_filemeta};
pub(crate) type ReplicationObjectStore = crate::store::ECStore;
pub(crate) use crate::client::api_get_options::{AdvancedGetOptions, StatObjectOptions};
pub(crate) use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
@@ -101,7 +102,7 @@ pub(crate) fn deleted_object_for_replication(delete_object: DeletedObject) -> Re
object_name: delete_object.object_name,
version_id: delete_object.version_id,
delete_marker_mtime: delete_object.delete_marker_mtime,
replication_state: delete_object.replication_state,
replication_state: delete_object.replication_state.as_ref().map(replication_state_from_filemeta),
found: delete_object.found,
force_delete: delete_object.force_delete,
}
@@ -112,7 +113,7 @@ pub(crate) fn object_to_delete_for_replication(object: &ObjectToDelete) -> Repli
object_name: object.object_name.clone(),
version_id: object.version_id,
delete_marker_replication_status: object.delete_marker_replication_status.clone(),
version_purge_status: object.version_purge_status.clone(),
version_purge_status: object.version_purge_status.clone().map(version_purge_status_from_filemeta),
version_purge_statuses: object.version_purge_statuses.clone(),
replicate_decision_str: object.replicate_decision_str.clone(),
}
@@ -21,7 +21,7 @@ const EVENT_LIFECYCLE_CLEANUP_SKIPPED: &str = "lifecycle_cleanup_skipped";
const EVENT_LIFECYCLE_CLEANUP_FAILED: &str = "lifecycle_cleanup_failed";
use crate::bucket::lifecycle::lifecycle;
use crate::bucket::replication::{ReplicationLifecycleBridge, ReplicationState};
use crate::bucket::replication::{ReplicationLifecycleBridge, ReplicationState, replication_state_to_filemeta};
use crate::bucket::versioning::VersioningApi;
use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::object_api::ObjectOptions;
@@ -106,7 +106,7 @@ pub async fn delete_object_versions(api: &Arc<ECStore>, bucket: &str, to_del: &[
let Some(replication_state) = replication_candidates.get(i).and_then(|c| c.clone()) else {
continue;
};
deleted_obj.replication_state = Some(replication_state);
deleted_obj.replication_state = Some(replication_state_to_filemeta(&replication_state));
ReplicationLifecycleBridge::schedule_delete(bucket.to_string(), deleted_obj.clone()).await;
}
+7 -3
View File
@@ -12,6 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::bucket::replication::replication_state_from_filemeta;
use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::bucket::{
lifecycle::{
@@ -2132,7 +2133,10 @@ fn decommission_delete_marker_opts(
data_movement: true,
delete_marker: true,
skip_decommissioned: true,
delete_replication: version.replication_state_internal.clone(),
delete_replication: version
.replication_state_internal
.as_ref()
.map(replication_state_from_filemeta),
..Default::default()
}
}
@@ -4236,12 +4240,12 @@ mod tests {
let mod_time = OffsetDateTime::now_utc();
let version = rustfs_filemeta::FileInfo {
mod_time: Some(mod_time),
replication_state_internal: Some(ReplicationState {
replication_state_internal: Some(crate::bucket::replication::replication_state_to_filemeta(&ReplicationState {
replica_status: ReplicationStatusType::Replica,
delete_marker: true,
replicate_decision_str: "existing".to_string(),
..Default::default()
}),
})),
..Default::default()
};
+2 -2
View File
@@ -17,8 +17,8 @@
use crate::bucket::metadata_sys::get_versioning_config;
use crate::bucket::replication::{
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_statuses_map,
version_purge_statuses_map,
ReplicateDecision, ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_status_from_filemeta,
replication_statuses_map, version_purge_status_from_filemeta, version_purge_statuses_map,
};
use crate::bucket::versioning::VersioningApi as _;
use crate::config::storageclass;
+4 -4
View File
@@ -445,7 +445,7 @@ impl ObjectInfo {
.map(|v| v.replicate_decision_str.clone())
.unwrap_or_default();
let mut replication_status = fi.replication_status();
let mut replication_status = replication_status_from_filemeta(fi.replication_status());
if replication_status.is_empty()
&& let Some(status) = fi.metadata.get(AMZ_BUCKET_REPLICATION_STATUS).cloned()
&& status == ReplicationStatusType::Replica.as_str()
@@ -453,7 +453,7 @@ impl ObjectInfo {
replication_status = ReplicationStatusType::Replica;
}
let version_purge_status = fi.version_purge_status();
let version_purge_status = version_purge_status_from_filemeta(fi.version_purge_status());
let transitioned_object = TransitionedObject {
name: fi.transitioned_objname.clone(),
@@ -1052,10 +1052,10 @@ mod tests {
#[test]
fn from_file_info_preserves_replication_decision() {
let fi = FileInfo {
replication_state_internal: Some(ReplicationState {
replication_state_internal: Some(crate::bucket::replication::replication_state_to_filemeta(&ReplicationState {
replicate_decision_str: "arn=true;false;arn:replication::1:dest;rule-id".to_string(),
..Default::default()
}),
})),
..Default::default()
};
@@ -1,4 +1,5 @@
use super::worker::{is_transient_rebalance_error, rebalance_migration_retry_delay, sleep_rebalance_migration_retry};
use crate::bucket::replication::replication_state_from_filemeta;
use crate::data_usage::DATA_USAGE_CACHE_NAME;
use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found};
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions};
@@ -32,7 +33,10 @@ pub(super) fn rebalance_delete_marker_opts(version: &FileInfo, version_id: Optio
data_movement: true,
delete_marker: true,
skip_decommissioned: true,
delete_replication: version.replication_state_internal.clone(),
delete_replication: version
.replication_state_internal
.as_ref()
.map(replication_state_from_filemeta),
..Default::default()
}
}
@@ -49,7 +49,7 @@ use super::{
DiskStat, GetObjectReader, ObjectInfo, ObjectOptions, RebalSaveOpt, RebalStatus, RebalanceBucketOutcome,
RebalanceCleanupWarnings, RebalanceEntryOutcome, RebalanceInfo, RebalanceMeta, RebalanceStats,
};
use crate::bucket::replication::{ReplicationState, ReplicationStatusType};
use crate::bucket::replication::{ReplicationState, ReplicationStatusType, replication_state_to_filemeta};
use crate::data_movement;
use crate::data_usage::DATA_USAGE_CACHE_NAME;
use crate::disk::RUSTFS_META_BUCKET;
@@ -223,12 +223,12 @@ fn test_rebalance_delete_marker_opts_preserves_replication_state() {
let mod_time = OffsetDateTime::now_utc();
let version = FileInfo {
mod_time: Some(mod_time),
replication_state_internal: Some(ReplicationState {
replication_state_internal: Some(replication_state_to_filemeta(&ReplicationState {
replica_status: ReplicationStatusType::Replica,
delete_marker: true,
replicate_decision_str: "existing".to_string(),
..Default::default()
}),
})),
..version_deleted()
};
+4 -3
View File
@@ -22,6 +22,7 @@ use crate::bucket::metadata_sys;
use crate::bucket::object_lock::objectlock_sys::check_retention_for_modification;
use crate::bucket::replication::{
ReplicateDecision, ReplicationObjectBridge, ReplicationState, ReplicationStatusType, VersionPurgeStatusType,
replication_state_to_filemeta,
};
use crate::bucket::versioning::VersioningApi;
use crate::bucket::versioning_sys::BucketVersioningSys;
@@ -2944,7 +2945,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
}
record_capacity_scope_if_needed(opts.capacity_scope_token, &online_disks);
fi.replication_state_internal = Some(opts.put_replication_state());
fi.replication_state_internal = Some(replication_state_to_filemeta(&opts.put_replication_state()));
fi.is_latest = true;
@@ -4152,7 +4153,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
deleted: delete_marker,
mark_deleted: mark_delete,
mod_time: Some(mod_time),
replication_state_internal: opts.delete_replication.clone(),
replication_state_internal: opts.delete_replication.as_ref().map(replication_state_to_filemeta),
..Default::default() // TODO: Transition
};
@@ -4190,7 +4191,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
mark_deleted: mark_delete,
deleted: delete_marker,
mod_time: Some(mod_time),
replication_state_internal: opts.delete_replication.clone(),
replication_state_internal: opts.delete_replication.as_ref().map(replication_state_to_filemeta),
..Default::default()
};
+4 -3
View File
@@ -1347,7 +1347,8 @@ mod tests {
use super::*;
use crate::bucket::lifecycle::core::TRANSITION_COMPLETE;
use crate::bucket::replication::{
ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_statuses_map, version_purge_statuses_map,
ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta, replication_statuses_map,
version_purge_statuses_map,
};
use crate::layout::{
endpoints::{Endpoints, PoolEndpoints},
@@ -1513,13 +1514,13 @@ mod tests {
transitioned_objname: "remote/object".to_string(),
transition_tier: "WARM".to_string(),
transition_version_id: Some(transition_version_id),
replication_state_internal: Some(ReplicationState {
replication_state_internal: Some(replication_state_to_filemeta(&ReplicationState {
replication_status_internal: Some("arn:minio:replication:target=COMPLETED;".to_string()),
targets: replication_statuses_map("arn:minio:replication:target=COMPLETED;"),
version_purge_status_internal: Some("arn:minio:replication:target=PENDING;".to_string()),
purge_targets: version_purge_statuses_map("arn:minio:replication:target=PENDING;"),
..Default::default()
}),
})),
metadata: HashMap::from([
("etag".to_string(), "etag-value".to_string()),
("x-amz-meta-key".to_string(), "metadata-value".to_string()),
+2 -1
View File
@@ -26,10 +26,11 @@ categories = ["web-programming", "development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-replication/latest/rustfs_replication/"
[dependencies]
bytes.workspace = true
byteorder.workspace = true
regex.workspace = true
rmp.workspace = true
rmp-serde.workspace = true
rustfs-filemeta.workspace = true
rustfs-utils = { workspace = true, features = ["http", "path", "string"] }
s3s.workspace = true
serde.workspace = true
File diff suppressed because it is too large Load Diff
+4 -3
View File
@@ -36,9 +36,10 @@ pub use delete::{
should_retry_delete_marker_purge,
};
pub use 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,
REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING,
REPLICATE_INCOMING_DELETE, REPLICATE_MRF, REPLICATE_QUEUED, REPLICATION_RESET, REPLICATION_STATUS, 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,
};