diff --git a/crates/ecstore/src/bucket/replication/mod.rs b/crates/ecstore/src/bucket/replication/mod.rs index 2d5087d8e..c939b3d41 100644 --- a/crates/ecstore/src/bucket/replication/mod.rs +++ b/crates/ecstore/src/bucket/replication/mod.rs @@ -18,6 +18,7 @@ mod replication_event_sink; mod replication_pool; mod replication_resyncer; mod replication_state; +mod replication_target_boundary; mod rule; mod runtime_boundary; diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 78a55bfa3..390da987e 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -12,8 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +use super::replication_target_boundary as target_boundary; use super::runtime_boundary as runtime_sources; -use crate::bucket::bucket_target_sys::BucketTargetSys; use crate::bucket::metadata_sys; use crate::bucket::replication::ResyncOpts; use crate::bucket::replication::ResyncStatusType; @@ -1508,7 +1508,7 @@ pub async fn queue_replication_heal(bucket: &str, oi: ObjectInfo, retry_count: u } }; - let tgts = match BucketTargetSys::get().list_bucket_targets(bucket).await { + let tgts = match target_boundary::list_bucket_targets(bucket).await { Ok(targets) => Some(targets), Err(err) => { debug!( diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index d9512be09..c6f4081bb 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -13,10 +13,11 @@ // limitations under the License. use super::replication_event_sink::{EventArgs, send_event}; +use super::replication_target_boundary as target_boundary; use super::runtime_boundary as runtime_sources; use crate::bucket::bandwidth::reader::{BucketOptions, MonitorReaderOptions, MonitoredReader}; use crate::bucket::bucket_target_sys::{ - AdvancedPutOptions, BucketTargetSys, PutObjectOptions, PutObjectPartOptions, RemoveObjectOptions, TargetClient, + AdvancedPutOptions, PutObjectOptions, PutObjectPartOptions, RemoveObjectOptions, TargetClient, }; use crate::bucket::metadata_sys; use crate::bucket::msgp_decode::{read_msgp_ext8_time, skip_msgp_value, write_msgp_time}; @@ -874,7 +875,7 @@ impl ReplicationResyncer { } }; - let targets = match BucketTargetSys::get().list_bucket_targets(&opts.bucket).await { + let targets = match target_boundary::list_bucket_targets(&opts.bucket).await { Ok(targets) => targets, Err(err) => { debug!( @@ -919,10 +920,7 @@ impl ReplicationResyncer { return; } - let Some(target_client) = BucketTargetSys::get() - .get_remote_target_client(&opts.bucket, &target_arns[0]) - .await - else { + let Some(target_client) = target_boundary::remote_target_client(&opts.bucket, &target_arns[0]).await else { error!( event = EVENT_RESYNC_RUNTIME_SKIPPED, component = LOG_COMPONENT_ECSTORE, @@ -1686,7 +1684,7 @@ pub async fn check_replicate_delete( continue; } - let tgt = BucketTargetSys::get().get_remote_target_client(bucket, &tgt_arn).await; + let tgt = target_boundary::remote_target_client(bucket, &tgt_arn).await; // The target online status should not be used here while deciding // whether to replicate deletes as the target could be temporarily down let tgt_dsc = if let Some(tgt) = tgt { @@ -1811,7 +1809,7 @@ pub async fn must_replicate(bucket: &str, object: &str, mopts: MustReplicateOpti let mut dsc = ReplicateDecision::default(); for arn in arns { - let cli = BucketTargetSys::get().get_remote_target_client(bucket, &arn).await; + let cli = target_boundary::remote_target_client(bucket, &arn).await; let mut sopts = opts.clone(); sopts.target_arn = arn.clone(); @@ -2078,7 +2076,7 @@ pub async fn replicate_delete(dobj: DeletedObjectReplicat } // Get the remote target client - let Some(tgt_client) = BucketTargetSys::get().get_remote_target_client(&bucket, &tgt_entry.arn).await else { + let Some(tgt_client) = target_boundary::remote_target_client(&bucket, &tgt_entry.arn).await else { debug!( event = EVENT_REPLICATION_DELETE_SKIPPED, component = LOG_COMPONENT_ECSTORE, @@ -2304,7 +2302,7 @@ async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedOb if !dobj.target_arn.is_empty() && dobj.target_arn != tgt_entry.arn { continue; } - let Some(tgt_client) = BucketTargetSys::get().get_remote_target_client(bucket, &tgt_entry.arn).await else { + let Some(tgt_client) = target_boundary::remote_target_client(bucket, &tgt_entry.arn).await else { continue; }; @@ -2455,7 +2453,7 @@ async fn replicate_force_delete_to_targets(dobj: &Deleted let mut join_set = JoinSet::new(); for arn in tgt_arns { - let Some(tgt_client) = BucketTargetSys::get().get_remote_target_client(bucket, &arn).await else { + let Some(tgt_client) = target_boundary::remote_target_client(bucket, &arn).await else { debug!( event = EVENT_REPLICATION_FORCE_DELETE_SKIPPED, component = LOG_COMPONENT_ECSTORE, @@ -2484,7 +2482,7 @@ async fn replicate_force_delete_to_targets(dobj: &Deleted let object_name = object_name.clone(); join_set.spawn(async move { - if BucketTargetSys::get().is_offline(&tgt_client.to_url()).await { + if target_boundary::target_is_offline(&tgt_client).await { error!( event = EVENT_REPLICATION_FORCE_DELETE_SKIPPED, component = LOG_COMPONENT_ECSTORE, @@ -2612,7 +2610,7 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli return rinfo; } - if BucketTargetSys::get().is_offline(&tgt_client.to_url()).await { + if target_boundary::target_is_offline(&tgt_client).await { if !is_version_purge { rinfo.replication_status = ReplicationStatusType::Failed; } else { @@ -2841,7 +2839,7 @@ pub async fn replicate_object(roi: ReplicateObjectInfo, s let mut join_set = JoinSet::new(); for arn in tgt_arns { - let Some(tgt_client) = BucketTargetSys::get().get_remote_target_client(&bucket, &arn).await else { + let Some(tgt_client) = target_boundary::remote_target_client(&bucket, &arn).await else { debug!( event = EVENT_RESYNC_RUNTIME_SKIPPED, component = LOG_COMPONENT_ECSTORE, @@ -2999,7 +2997,7 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo { return rinfo; } - if BucketTargetSys::get().is_offline(&tgt_client.to_url()).await { + if target_boundary::target_is_offline(&tgt_client).await { debug!( event = EVENT_RESYNC_RUNTIME_SKIPPED, component = LOG_COMPONENT_ECSTORE, @@ -3286,7 +3284,7 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo { ..Default::default() }; - if BucketTargetSys::get().is_offline(&tgt_client.to_url()).await { + if target_boundary::target_is_offline(&tgt_client).await { debug!( event = EVENT_RESYNC_RUNTIME_SKIPPED, component = LOG_COMPONENT_ECSTORE, diff --git a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs new file mode 100644 index 000000000..20906a60b --- /dev/null +++ b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs @@ -0,0 +1,30 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use crate::bucket::bucket_target_sys::{BucketTargetError, BucketTargetSys, TargetClient}; +use crate::bucket::target::BucketTargets; + +pub(crate) async fn list_bucket_targets(bucket: &str) -> Result { + BucketTargetSys::get().list_bucket_targets(bucket).await +} + +pub(crate) async fn remote_target_client(bucket: &str, arn: &str) -> Option> { + BucketTargetSys::get().get_remote_target_client(bucket, arn).await +} + +pub(crate) async fn target_is_offline(target_client: &TargetClient) -> bool { + BucketTargetSys::get().is_offline(&target_client.to_url()).await +} diff --git a/scripts/check_architecture_migration_rules.sh b/scripts/check_architecture_migration_rules.sh index 0ebdfd722..b982497db 100755 --- a/scripts/check_architecture_migration_rules.sh +++ b/scripts/check_architecture_migration_rules.sh @@ -184,6 +184,7 @@ FUZZ_ECSTORE_COMPAT_BYPASS_HITS_FILE="${TMP_DIR}/fuzz_ecstore_compat_bypass_hits EXTERNAL_ECSTORE_API_BOUNDARY_HITS_FILE="${TMP_DIR}/external_ecstore_api_boundary_hits.txt" REPLICATION_FACADE_BYPASS_HITS_FILE="${TMP_DIR}/replication_facade_bypass_hits.txt" REPLICATION_EVENT_SINK_BYPASS_HITS_FILE="${TMP_DIR}/replication_event_sink_bypass_hits.txt" +REPLICATION_TARGET_BOUNDARY_BYPASS_HITS_FILE="${TMP_DIR}/replication_target_boundary_bypass_hits.txt" REPLICATION_RUNTIME_SOURCE_BYPASS_HITS_FILE="${TMP_DIR}/replication_runtime_source_bypass_hits.txt" GLOBAL_REPLICATION_STATE_BYPASS_HITS_FILE="${TMP_DIR}/global_replication_state_bypass_hits.txt" GLOBAL_BUCKET_MONITOR_BYPASS_HITS_FILE="${TMP_DIR}/global_bucket_monitor_bypass_hits.txt" @@ -2351,6 +2352,18 @@ if [[ -s "$REPLICATION_EVENT_SINK_BYPASS_HITS_FILE" ]]; then report_failure "replication event notification access must stay behind replication event sink: $(paste -sd '; ' "$REPLICATION_EVENT_SINK_BYPASS_HITS_FILE")" fi +( + cd "$ROOT_DIR" + rg -n --with-filename 'BucketTargetSys::get\(\)' \ + crates/ecstore/src/bucket/replication \ + --glob '*.rs' | + rg -v '^crates/ecstore/src/bucket/replication/replication_target_boundary\.rs:' || true +) >"$REPLICATION_TARGET_BOUNDARY_BYPASS_HITS_FILE" + +if [[ -s "$REPLICATION_TARGET_BOUNDARY_BYPASS_HITS_FILE" ]]; then + report_failure "replication bucket-target access must stay behind replication target boundary: $(paste -sd '; ' "$REPLICATION_TARGET_BOUNDARY_BYPASS_HITS_FILE")" +fi + ( cd "$ROOT_DIR" rg -n --with-filename '\bGLOBAL_REPLICATION_(POOL|STATS)\b' \