mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-09 13:46:05 +00:00
fix: preserve scanner backlog during pool retirement (#7555)
* fix: hand off native scanner backlog before pool retirement * test: initialize optional proof in data usage fixtures * perf: bound parallel V3 pool metadata phase writes * test: keep native backlog fault errors typed * test: skip unchanged pool metadata in crash barriers --------- Co-authored-by: cxymds <cxymds@gmail.com>
This commit is contained in:
@@ -367,6 +367,14 @@ pub mod config {
|
||||
}
|
||||
|
||||
pub mod data_usage {
|
||||
#[cfg(feature = "test-util")]
|
||||
pub use crate::data_movement::SourceCleanupDeleteBarrier;
|
||||
#[cfg(feature = "test-util")]
|
||||
pub use crate::data_movement::scanner_backlog::test_util::NativeScannerPauseBacklogWriteFault;
|
||||
pub use crate::data_movement::scanner_backlog::{
|
||||
MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementPlanner,
|
||||
ScannerPauseBacklogRetirementReplica, register_scanner_pause_backlog_retirement_planner,
|
||||
};
|
||||
pub use crate::data_usage::{
|
||||
DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage,
|
||||
init_compression_total_memory_from_backend, invalidate_admin_data_usage_snapshot_cache,
|
||||
|
||||
+1788
-62
File diff suppressed because it is too large
Load Diff
@@ -5256,16 +5256,35 @@ mod decommission_lock_order_tests {
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial]
|
||||
fn scanner_backlog_native_replica_reconciles_capacity_and_cleans_source() {
|
||||
run_large_stack_current_thread_async_test("scanner-backlog-reconcile", async || {
|
||||
fn data_movement_existing_replica_reconciles_capacity_and_cleans_source() {
|
||||
data_movement_existing_replica_reconciles_capacity_case(false);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial]
|
||||
fn data_movement_existing_replica_outside_reservation_uses_reserved_target() {
|
||||
data_movement_existing_replica_reconciles_capacity_case(true);
|
||||
}
|
||||
|
||||
fn data_movement_existing_replica_reconciles_capacity_case(existing_outside_reservation: bool) {
|
||||
run_large_stack_current_thread_async_test("reserved-replica-reconcile", async move || {
|
||||
let (_temp_dirs, store, other_store) =
|
||||
test_three_pool_stores_with_three_disk_sets_with_isolated_node_contexts(None).await;
|
||||
let object = "buckets/.scanner-pause-backlog.json";
|
||||
let object = "buckets/reserved-replica-routing.json";
|
||||
let body = br#"{"schemaVersion":1,"generation":2}"#.to_vec();
|
||||
let old_body = br#"{"schemaVersion":1,"generation":1}"#.to_vec();
|
||||
let source_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(20);
|
||||
let target_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(10);
|
||||
for (pool_index, payload, mod_time) in [(0, body.clone(), source_time), (2, old_body, target_time)] {
|
||||
let target_time = source_time;
|
||||
let target_pool_index = if existing_outside_reservation { 1 } else { 2 };
|
||||
let mut replicas = vec![(0, body.clone(), source_time), (target_pool_index, old_body, target_time)];
|
||||
if existing_outside_reservation {
|
||||
replicas.push((
|
||||
2,
|
||||
br#"{"schemaVersion":1,"generation":3}"#.to_vec(),
|
||||
source_time + time::Duration::seconds(10),
|
||||
));
|
||||
}
|
||||
for (pool_index, payload, mod_time) in replicas.iter().cloned() {
|
||||
store.pools[pool_index]
|
||||
.put_object(
|
||||
RUSTFS_META_BUCKET,
|
||||
@@ -5278,13 +5297,19 @@ mod decommission_lock_order_tests {
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("seed native scanner replicas with independent write times");
|
||||
.expect("seed existing replicas with independent write times");
|
||||
}
|
||||
let layout = DecommissionErasureLayout { data: 1, parity: 0 };
|
||||
let target_total = body.len() * 8;
|
||||
let capacities = vec![
|
||||
DecommissionPoolCapacityInfo::for_test(0, layout, 0, body.len() * 2, body.len() * 2),
|
||||
DecommissionPoolCapacityInfo::for_test(1, layout, 0, target_total, target_total),
|
||||
DecommissionPoolCapacityInfo::for_test(
|
||||
1,
|
||||
layout,
|
||||
if existing_outside_reservation { target_total } else { 0 },
|
||||
target_total,
|
||||
if existing_outside_reservation { 0 } else { target_total },
|
||||
),
|
||||
DecommissionPoolCapacityInfo::for_test(2, layout, target_total, target_total, 0),
|
||||
];
|
||||
set_decommission_capacity_info_overrides_for_test(store.id, vec![capacities.clone()]);
|
||||
@@ -5293,6 +5318,17 @@ mod decommission_lock_order_tests {
|
||||
.await
|
||||
.expect("activate the source reservation");
|
||||
let owner = decommission_capacity_owner(&*store.pool_meta.read().await);
|
||||
let reserved_snapshot = store.pool_meta.read().await.clone();
|
||||
let reservation = reserved_snapshot.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.and_then(|info| info.capacity_reservation.as_ref())
|
||||
.expect("active source reservation");
|
||||
assert_eq!(
|
||||
reservation.targets.iter().map(|target| target.pool_index).collect::<Vec<_>>(),
|
||||
vec![target_pool_index],
|
||||
"the fixture must reserve exactly one target"
|
||||
);
|
||||
let source_reader = store.pools[0]
|
||||
.get_object_reader(
|
||||
RUSTFS_META_BUCKET,
|
||||
@@ -5314,12 +5350,91 @@ mod decommission_lock_order_tests {
|
||||
RUSTFS_META_BUCKET.to_string(),
|
||||
source_reader,
|
||||
None,
|
||||
"scanner_backlog_conflict",
|
||||
"reserved_replica_conflict",
|
||||
Some(owner),
|
||||
)
|
||||
.await
|
||||
.expect_err("a different older native ledger must retain its source and capacity intent");
|
||||
.expect_err("a different older existing record must retain its source and capacity intent");
|
||||
assert!(conflict.to_string().contains("Precondition failed"), "unexpected conflict: {conflict}");
|
||||
let reserved_snapshot = store.pool_meta.read().await.clone();
|
||||
let mut selection_opts = ObjectOptions {
|
||||
data_movement: true,
|
||||
src_pool_idx: 0,
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(
|
||||
store
|
||||
.select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &selection_opts, true)
|
||||
.await
|
||||
.expect("selection without a capacity owner retains existing-replica routing"),
|
||||
2
|
||||
);
|
||||
owner.apply_to(&mut selection_opts);
|
||||
for stale_owner in [
|
||||
DecommissionCapacityOwner {
|
||||
owner_nonce: uuid::Uuid::new_v4(),
|
||||
..owner
|
||||
},
|
||||
DecommissionCapacityOwner {
|
||||
generation: owner.generation + 1,
|
||||
..owner
|
||||
},
|
||||
] {
|
||||
let mut stale_opts = selection_opts.clone();
|
||||
stale_owner.apply_to(&mut stale_opts);
|
||||
assert!(
|
||||
matches!(
|
||||
store
|
||||
.select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &stale_opts, true)
|
||||
.await,
|
||||
Err(crate::error::Error::DecommissionCapacityBlocked { .. })
|
||||
),
|
||||
"a stale owner must not fall back to another target"
|
||||
);
|
||||
}
|
||||
{
|
||||
let mut meta = store.pool_meta.write().await;
|
||||
meta.pools[0]
|
||||
.decommission
|
||||
.as_mut()
|
||||
.unwrap()
|
||||
.capacity_reservation
|
||||
.as_mut()
|
||||
.unwrap()
|
||||
.expires_at = time::OffsetDateTime::now_utc() - time::Duration::seconds(1);
|
||||
}
|
||||
assert!(
|
||||
matches!(
|
||||
store
|
||||
.select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &selection_opts, true)
|
||||
.await,
|
||||
Err(crate::error::Error::DecommissionCapacityBlocked { .. })
|
||||
),
|
||||
"an expired owner must not fall back to another target"
|
||||
);
|
||||
*store.pool_meta.write().await = reserved_snapshot.clone();
|
||||
if !existing_outside_reservation {
|
||||
{
|
||||
let mut meta = store.pool_meta.write().await;
|
||||
let target = &mut meta.pools[0]
|
||||
.decommission
|
||||
.as_mut()
|
||||
.unwrap()
|
||||
.capacity_reservation
|
||||
.as_mut()
|
||||
.unwrap()
|
||||
.targets[0];
|
||||
target.consumed_physical_bytes = target.reserved_physical_bytes;
|
||||
}
|
||||
assert_eq!(
|
||||
store
|
||||
.select_data_movement_pool_idx(RUSTFS_META_BUCKET, object, body.len() as i64, &selection_opts, true)
|
||||
.await
|
||||
.expect("an existing reserved replica can still be selected after capacity was consumed"),
|
||||
target_pool_index
|
||||
);
|
||||
*store.pool_meta.write().await = reserved_snapshot;
|
||||
}
|
||||
let mut persisted = crate::core::pools::PoolMeta::default();
|
||||
persisted
|
||||
.load_no_lock_from_replicas(store.pools.clone())
|
||||
@@ -5336,11 +5451,24 @@ mod decommission_lock_order_tests {
|
||||
.pending_target_physical_bytes,
|
||||
body.len()
|
||||
);
|
||||
let previous = store.pools[2]
|
||||
for (pool_index, payload, mod_time) in &replicas {
|
||||
let mut reader = store.pools[*pool_index]
|
||||
.get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("a refused existing record replacement must preserve every replica");
|
||||
assert_eq!(reader.object_info.mod_time, Some(*mod_time));
|
||||
let mut actual = Vec::new();
|
||||
reader
|
||||
.read_to_end(&mut actual)
|
||||
.await
|
||||
.expect("read the unchanged existing record");
|
||||
assert_eq!(&actual, payload);
|
||||
}
|
||||
let previous = store.pools[target_pool_index]
|
||||
.get_object_info(RUSTFS_META_BUCKET, object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("read the native writer's CAS revision");
|
||||
let replacement = store.pools[2]
|
||||
.expect("read the existing writer's CAS revision");
|
||||
let replacement = store.pools[target_pool_index]
|
||||
.put_object(
|
||||
RUSTFS_META_BUCKET,
|
||||
object,
|
||||
@@ -5356,7 +5484,7 @@ mod decommission_lock_order_tests {
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("native scanner CAS converges the payload without a migration marker");
|
||||
.expect("existing CAS converges the payload without a migration marker");
|
||||
assert!(!data_movement::is_owned_data_movement_target(&replacement));
|
||||
*other_store.pool_meta.write().await = persisted;
|
||||
set_decommission_capacity_info_overrides_for_test(other_store.id, vec![capacities]);
|
||||
@@ -5374,7 +5502,7 @@ mod decommission_lock_order_tests {
|
||||
)
|
||||
.await
|
||||
.expect("replica conflict recovery must be bounded")
|
||||
.expect("identical native replica should finish migration on the reloaded node");
|
||||
.expect("identical existing replica should finish migration on the reloaded node");
|
||||
let mut reconciled = crate::core::pools::PoolMeta::default();
|
||||
reconciled
|
||||
.load_no_lock_from_replicas(other_store.pools.clone())
|
||||
@@ -5404,14 +5532,14 @@ mod decommission_lock_order_tests {
|
||||
.await
|
||||
.expect_err("the source should be cleaned only after equivalent-target capacity reconciliation");
|
||||
assert!(crate::error::is_err_object_not_found(&missing));
|
||||
let mut target_reader = other_store.pools[2]
|
||||
let mut target_reader = other_store.pools[target_pool_index]
|
||||
.get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("the surviving replica should remain readable");
|
||||
assert_eq!(
|
||||
target_reader.object_info.mod_time,
|
||||
Some(target_time),
|
||||
"recovery must not overwrite the native target"
|
||||
"recovery must not overwrite the existing target"
|
||||
);
|
||||
let mut actual = Vec::new();
|
||||
target_reader
|
||||
@@ -5419,6 +5547,20 @@ mod decommission_lock_order_tests {
|
||||
.await
|
||||
.expect("read surviving ledger bytes");
|
||||
assert_eq!(actual, body);
|
||||
if existing_outside_reservation {
|
||||
let (_, outside_body, outside_time) = replicas.last().expect("unreserved existing replica");
|
||||
let mut outside = other_store.pools[2]
|
||||
.get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("migration must leave the unreserved existing replica intact");
|
||||
assert_eq!(outside.object_info.mod_time, Some(*outside_time));
|
||||
let mut actual = Vec::new();
|
||||
outside
|
||||
.read_to_end(&mut actual)
|
||||
.await
|
||||
.expect("read the untouched unreserved replica");
|
||||
assert_eq!(&actual, outside_body);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -15,6 +15,7 @@
|
||||
// #730: data-movement migration keeps staged cleanup helpers until copy paths converge.
|
||||
|
||||
pub(crate) mod backpressure;
|
||||
pub(crate) mod scanner_backlog;
|
||||
|
||||
use crate::core::pools::{DecommissionCapacityOwner, decommission_capacity_mutation_id};
|
||||
use crate::error::{
|
||||
@@ -984,24 +985,6 @@ fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target:
|
||||
.is_some_and(|(source_time, target_time)| target_time > source_time)
|
||||
}
|
||||
|
||||
fn is_equivalent_scanner_backlog_replica(source: &ObjectInfo, target: &ObjectInfo, compare_part_checksums: bool) -> bool {
|
||||
// Scanner publishes this exact payload to surviving sets with CAS. Each
|
||||
// set assigns its own write time; that timestamp is not a ledger generation.
|
||||
// Accept only an identical, known unversioned identity, never a different
|
||||
// record based on timestamp ordering or a similarly named user object.
|
||||
source.bucket == crate::disk::RUSTFS_META_BUCKET
|
||||
&& target.bucket == source.bucket
|
||||
&& source.name == "buckets/.scanner-pause-backlog.json"
|
||||
&& target.name == source.name
|
||||
&& is_unversioned_data_movement_object(source)
|
||||
&& is_unversioned_data_movement_object(target)
|
||||
&& !source.delete_marker
|
||||
&& source.mod_time.is_some()
|
||||
&& target.mod_time.is_some()
|
||||
&& source.etag.as_ref().is_some_and(|etag| !etag.is_empty())
|
||||
&& is_equivalent_data_movement_object_identity(source, target, false, compare_part_checksums)
|
||||
}
|
||||
|
||||
fn is_data_movement_upload_takeover_target(source: &ObjectInfo, target: &ObjectInfo, compare_part_checksums: bool) -> bool {
|
||||
let identity = data_movement_upload_identity(source);
|
||||
source.mod_time.is_some()
|
||||
@@ -1217,7 +1200,7 @@ struct SourceCleanupDeleteBarrierState {
|
||||
dead_code,
|
||||
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
|
||||
)]
|
||||
pub(crate) struct SourceCleanupDeleteBarrier {
|
||||
pub struct SourceCleanupDeleteBarrier {
|
||||
state: Arc<SourceCleanupDeleteBarrierState>,
|
||||
}
|
||||
|
||||
@@ -1231,7 +1214,7 @@ static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<
|
||||
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
|
||||
)]
|
||||
impl SourceCleanupDeleteBarrier {
|
||||
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
||||
pub fn install(bucket: &str, object: &str) -> Self {
|
||||
let state = Arc::new(SourceCleanupDeleteBarrierState {
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
@@ -1254,7 +1237,7 @@ impl SourceCleanupDeleteBarrier {
|
||||
Self { state }
|
||||
}
|
||||
|
||||
pub(crate) async fn wait_until_paused(&self) {
|
||||
pub async fn wait_until_paused(&self) {
|
||||
tokio::time::timeout(StdDuration::from_secs(30), self.state.arrived.notified())
|
||||
.await
|
||||
.expect("source cleanup should reach the pre-delete barrier");
|
||||
@@ -1270,7 +1253,7 @@ impl SourceCleanupDeleteBarrier {
|
||||
self.state.is_paused.load(Ordering::Acquire)
|
||||
}
|
||||
|
||||
pub(crate) fn release(&self) {
|
||||
pub fn release(&self) {
|
||||
self.state.release.notify_one();
|
||||
}
|
||||
}
|
||||
@@ -1449,7 +1432,8 @@ fn resolve_data_movement_overwrite_resume_result_for(
|
||||
target_pool_idx: usize,
|
||||
compare_part_checksums: bool,
|
||||
) -> Result<bool> {
|
||||
if !should_check_data_movement_overwrite_resume(err)
|
||||
if scanner_backlog::is_scanner_pause_backlog(&source.bucket, &source.name)
|
||||
|| !should_check_data_movement_overwrite_resume(err)
|
||||
|| !should_check_data_movement_resume_target(src_pool_idx, target_pool_idx)
|
||||
{
|
||||
return Ok(false);
|
||||
@@ -1471,9 +1455,7 @@ fn resolve_data_movement_overwrite_resume_result_for(
|
||||
return Ok(true);
|
||||
}
|
||||
|
||||
Ok(matches!(err, Error::PreconditionFailed)
|
||||
&& (is_equivalent_scanner_backlog_replica(source, &target, compare_part_checksums)
|
||||
|| is_superseding_unversioned_data_movement_object(source, &target)))
|
||||
Ok(matches!(err, Error::PreconditionFailed) && is_superseding_unversioned_data_movement_object(source, &target))
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
@@ -1646,6 +1628,9 @@ async fn migrate_object_inner(
|
||||
capacity_owner: Option<DecommissionCapacityOwner>,
|
||||
mutation_fence: Option<DecommissionFixedReadAnchor>,
|
||||
) -> Result<()> {
|
||||
if scanner_backlog::is_scanner_pause_backlog(&bucket, &rd.object_info.name) {
|
||||
return Err(Error::other("scanner pause backlog requires native retirement handoff"));
|
||||
}
|
||||
let mut mutation_fence = mutation_fence;
|
||||
let object_info = rd.object_info.clone();
|
||||
let capacity_owner = capacity_owner.map(|owner| {
|
||||
@@ -3354,16 +3339,25 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_scanner_backlog_resume_accepts_identical_native_replica_with_older_write_time() {
|
||||
fn test_scanner_backlog_resume_requires_native_cohort_proof_even_for_identical_payload() {
|
||||
let (source, target) = scanner_backlog_replica_pair();
|
||||
assert!(!is_owned_data_movement_target(&target), "native scanner writes are not migration copies");
|
||||
assert!(!is_equivalent_data_movement_object(&source, &target));
|
||||
assert!(
|
||||
scanner_backlog_precondition_resumes(&source, target),
|
||||
"identical ledger payloads have replica-local write times, not distinct committed generations"
|
||||
!scanner_backlog_precondition_resumes(&source, target),
|
||||
"a single identical replica cannot prove native cohort authority"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_scanner_backlog_resume_rejects_newer_timestamp_and_full_single_replica_identity() {
|
||||
let (source, mut target) = scanner_backlog_replica_pair();
|
||||
target.mod_time = source.mod_time.map(|time| time + time::Duration::SECOND);
|
||||
target.etag = Some("different-native-ledger".to_string());
|
||||
assert!(!scanner_backlog_precondition_resumes(&source, target));
|
||||
assert!(!scanner_backlog_precondition_resumes(&source, source.clone()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_scanner_backlog_resume_rejects_changed_payload_or_metadata() {
|
||||
let (source, target) = scanner_backlog_replica_pair();
|
||||
|
||||
@@ -0,0 +1,292 @@
|
||||
// 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 crate::disk::RUSTFS_META_BUCKET;
|
||||
use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found};
|
||||
use crate::object_api::ObjectOptions;
|
||||
use crate::object_api::{ObjectInfo, PutObjReader, WriteCompletion};
|
||||
use crate::set_disk::SetDisks;
|
||||
use crate::storage_api_contracts::object::HTTPPreconditions;
|
||||
use crate::storage_api_contracts::object::ObjectIO as _;
|
||||
use futures::future::join_all;
|
||||
use http::HeaderMap;
|
||||
use std::sync::{Arc, OnceLock};
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
pub const MAX_SCANNER_PAUSE_BACKLOG_BYTES: u64 = 64 * 1024;
|
||||
pub(crate) const SCANNER_PAUSE_BACKLOG_PATH: &str = "buckets/.scanner-pause-backlog.json";
|
||||
|
||||
/// A bounded, storage-fenced native replica. Only a confirmed missing object
|
||||
/// has no payload; read failures never enter the Scanner verifier.
|
||||
pub struct ScannerPauseBacklogRetirementReplica {
|
||||
pub pool_index: usize,
|
||||
pub set_index: usize,
|
||||
pub data: Option<Vec<u8>>,
|
||||
}
|
||||
|
||||
/// Native records for a membership handoff. Existing durable ledgers are
|
||||
/// preserved; an empty native bootstrap may initialize its first ledger.
|
||||
pub struct ScannerPauseBacklogRetirementPlan {
|
||||
pub seed_record: Option<Vec<u8>>,
|
||||
pub commit_record: Vec<u8>,
|
||||
pub stable_record: Vec<u8>,
|
||||
}
|
||||
|
||||
pub type ScannerPauseBacklogRetirementPlanner =
|
||||
fn(usize, &[ScannerPauseBacklogRetirementReplica]) -> std::result::Result<Option<ScannerPauseBacklogRetirementPlan>, String>;
|
||||
|
||||
static RETIREMENT_PLANNER: OnceLock<ScannerPauseBacklogRetirementPlanner> = OnceLock::new();
|
||||
|
||||
/// Install the stateless native record planner before storage starts workers.
|
||||
/// The scanner runtime switch does not control this storage safety check.
|
||||
pub fn register_scanner_pause_backlog_retirement_planner(planner: ScannerPauseBacklogRetirementPlanner) {
|
||||
RETIREMENT_PLANNER.get_or_init(|| planner);
|
||||
}
|
||||
|
||||
pub(crate) fn is_scanner_pause_backlog(bucket: &str, object: &str) -> bool {
|
||||
bucket == RUSTFS_META_BUCKET && object == SCANNER_PAUSE_BACKLOG_PATH
|
||||
}
|
||||
|
||||
pub(crate) struct ScannerPauseBacklogRetirementRead {
|
||||
pub replica: ScannerPauseBacklogRetirementReplica,
|
||||
pub etag: Option<String>,
|
||||
}
|
||||
|
||||
impl ScannerPauseBacklogRetirementRead {
|
||||
pub(crate) fn preconditions(&self) -> HTTPPreconditions {
|
||||
match &self.etag {
|
||||
Some(etag) => HTTPPreconditions {
|
||||
if_match: Some(etag.clone()),
|
||||
..Default::default()
|
||||
},
|
||||
None => HTTPPreconditions {
|
||||
if_none_match: Some("*".to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn read_replica(set: Arc<SetDisks>) -> Result<ScannerPauseBacklogRetirementRead> {
|
||||
let mut replica = ScannerPauseBacklogRetirementReplica {
|
||||
pool_index: set.pool_index,
|
||||
set_index: set.set_index,
|
||||
data: None,
|
||||
};
|
||||
let reader = match set
|
||||
.get_object_reader(
|
||||
RUSTFS_META_BUCKET,
|
||||
SCANNER_PAUSE_BACKLOG_PATH,
|
||||
None,
|
||||
HeaderMap::new(),
|
||||
&ObjectOptions {
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(reader) => reader,
|
||||
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {
|
||||
return Ok(ScannerPauseBacklogRetirementRead { replica, etag: None });
|
||||
}
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
let info = &reader.object_info;
|
||||
if info.version_id.is_some_and(|version| !version.is_nil())
|
||||
|| info.delete_marker
|
||||
|| info.is_dir
|
||||
|| info.etag.as_ref().is_none_or(String::is_empty)
|
||||
|| info.size < 0
|
||||
|| info.size > MAX_SCANNER_PAUSE_BACKLOG_BYTES as i64
|
||||
{
|
||||
return Err(Error::other("scanner pause backlog retirement found an unsupported replica identity"));
|
||||
}
|
||||
let etag = info.etag.clone();
|
||||
let expected_size = info.size as usize;
|
||||
let mut data = Vec::new();
|
||||
reader
|
||||
.take(MAX_SCANNER_PAUSE_BACKLOG_BYTES + 1)
|
||||
.read_to_end(&mut data)
|
||||
.await?;
|
||||
if data.len() != expected_size || data.len() > MAX_SCANNER_PAUSE_BACKLOG_BYTES as usize {
|
||||
return Err(Error::other("scanner pause backlog retirement replica has an invalid payload length"));
|
||||
}
|
||||
replica.data = Some(data);
|
||||
Ok(ScannerPauseBacklogRetirementRead { replica, etag })
|
||||
}
|
||||
|
||||
/// The caller retains the fixed object write lock and durable topology read
|
||||
/// fence through both this snapshot and physical source cleanup.
|
||||
pub(crate) async fn read_scanner_pause_backlog_retirement_replicas(
|
||||
source_pool_index: usize,
|
||||
source_set_index: usize,
|
||||
sets: Vec<Arc<SetDisks>>,
|
||||
) -> Result<Vec<ScannerPauseBacklogRetirementRead>> {
|
||||
let replicas = join_all(sets.into_iter().map(read_replica))
|
||||
.await
|
||||
.into_iter()
|
||||
.collect::<Result<Vec<_>>>()?;
|
||||
if !replicas.iter().any(|read| {
|
||||
read.replica.pool_index == source_pool_index && read.replica.set_index == source_set_index && read.replica.data.is_some()
|
||||
}) {
|
||||
return Err(Error::other("scanner pause backlog retirement current source replica is missing"));
|
||||
}
|
||||
Ok(replicas)
|
||||
}
|
||||
|
||||
pub(crate) fn plan_scanner_pause_backlog_retirement(
|
||||
source_pool_index: usize,
|
||||
replicas: &[ScannerPauseBacklogRetirementRead],
|
||||
) -> Result<Option<ScannerPauseBacklogRetirementPlan>> {
|
||||
let planner = RETIREMENT_PLANNER
|
||||
.get()
|
||||
.ok_or_else(|| Error::other("scanner pause backlog native retirement planner is unavailable"))?;
|
||||
let snapshots = replicas
|
||||
.iter()
|
||||
.map(|read| ScannerPauseBacklogRetirementReplica {
|
||||
pool_index: read.replica.pool_index,
|
||||
set_index: read.replica.set_index,
|
||||
data: read.replica.data.clone(),
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
planner(source_pool_index, &snapshots).map_err(Error::other)
|
||||
}
|
||||
|
||||
/// The native writer and retirement handoff use the same conditional, full-tail
|
||||
/// write. Their callers retain object and durable membership fences until return.
|
||||
pub(crate) async fn persist_native_scanner_pause_backlog_replica(
|
||||
set: Arc<SetDisks>,
|
||||
data: Vec<u8>,
|
||||
preconditions: HTTPPreconditions,
|
||||
mut opts: ObjectOptions,
|
||||
_phase: &'static str,
|
||||
) -> Result<ObjectInfo> {
|
||||
if data.len() > MAX_SCANNER_PAUSE_BACKLOG_BYTES as usize {
|
||||
return Err(Error::other("scanner pause backlog exceeds its size bound"));
|
||||
}
|
||||
opts.max_parity = true;
|
||||
opts.write_completion = WriteCompletion::TailDrained;
|
||||
opts.http_preconditions = Some(preconditions);
|
||||
#[cfg(feature = "test-util")]
|
||||
let fault = test_util::matching_write(&set, _phase)?;
|
||||
let result = set
|
||||
.put_object(RUSTFS_META_BUCKET, SCANNER_PAUSE_BACKLOG_PATH, &mut PutObjReader::from_vec(data), &opts)
|
||||
.await;
|
||||
#[cfg(feature = "test-util")]
|
||||
if result.is_ok()
|
||||
&& let Some(fault) = fault
|
||||
{
|
||||
fault.arrived.notify_one();
|
||||
fault.release.notified().await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
pub mod test_util {
|
||||
use super::*;
|
||||
use std::sync::Mutex;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use tokio::sync::Notify;
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
#[error("injected native scanner backlog {phase} write failure")]
|
||||
struct InjectedWriteFailure {
|
||||
phase: &'static str,
|
||||
}
|
||||
|
||||
pub(super) struct WriteFault {
|
||||
set: Arc<SetDisks>,
|
||||
phase: &'static str,
|
||||
remaining: AtomicUsize,
|
||||
fail_before_write: bool,
|
||||
pub(super) arrived: Notify,
|
||||
pub(super) release: Notify,
|
||||
}
|
||||
|
||||
static WRITE_FAULTS: Mutex<Vec<Arc<WriteFault>>> = Mutex::new(Vec::new());
|
||||
|
||||
/// Scope a one-shot fault to the actual set instance, so other stores and
|
||||
/// concurrent tests keep using the ordinary native persistence path.
|
||||
pub struct NativeScannerPauseBacklogWriteFault {
|
||||
state: Arc<WriteFault>,
|
||||
}
|
||||
|
||||
impl NativeScannerPauseBacklogWriteFault {
|
||||
fn install(set: Arc<SetDisks>, phase: &'static str, nth: usize, fail_before_write: bool) -> Self {
|
||||
assert!(nth > 0);
|
||||
let state = Arc::new(WriteFault {
|
||||
set,
|
||||
phase,
|
||||
remaining: AtomicUsize::new(nth),
|
||||
fail_before_write,
|
||||
arrived: Notify::new(),
|
||||
release: Notify::new(),
|
||||
});
|
||||
let mut faults = WRITE_FAULTS.lock().unwrap();
|
||||
assert!(
|
||||
!faults
|
||||
.iter()
|
||||
.any(|fault| Arc::ptr_eq(&fault.set, &state.set) && fault.phase == phase)
|
||||
);
|
||||
faults.push(Arc::clone(&state));
|
||||
Self { state }
|
||||
}
|
||||
|
||||
pub fn fail_before_write(set: Arc<SetDisks>, phase: &'static str, nth: usize) -> Self {
|
||||
Self::install(set, phase, nth, true)
|
||||
}
|
||||
|
||||
pub fn pause_after_write(set: Arc<SetDisks>, phase: &'static str) -> Self {
|
||||
Self::install(set, phase, 1, false)
|
||||
}
|
||||
|
||||
pub async fn wait_until_paused(&self) {
|
||||
self.state.arrived.notified().await;
|
||||
}
|
||||
|
||||
pub fn release(&self) {
|
||||
self.state.release.notify_one();
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for NativeScannerPauseBacklogWriteFault {
|
||||
fn drop(&mut self) {
|
||||
self.release();
|
||||
WRITE_FAULTS.lock().unwrap().retain(|fault| !Arc::ptr_eq(fault, &self.state));
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn matching_write(set: &Arc<SetDisks>, phase: &'static str) -> Result<Option<Arc<WriteFault>>> {
|
||||
let fault = WRITE_FAULTS
|
||||
.lock()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.find(|fault| Arc::ptr_eq(&fault.set, set) && fault.phase == phase)
|
||||
.cloned();
|
||||
let Some(fault) = fault else { return Ok(None) };
|
||||
if fault
|
||||
.remaining
|
||||
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |remaining| remaining.checked_sub(1))
|
||||
!= Ok(1)
|
||||
{
|
||||
return Ok(None);
|
||||
}
|
||||
if fault.fail_before_write {
|
||||
return Err(Error::other(InjectedWriteFailure { phase }));
|
||||
}
|
||||
Ok(Some(fault))
|
||||
}
|
||||
}
|
||||
@@ -3079,12 +3079,7 @@ impl ECStore {
|
||||
let store = Arc::clone(self);
|
||||
let write = async move {
|
||||
let object = "buckets/.scanner-pause-backlog.json";
|
||||
let mut opts = ObjectOptions {
|
||||
max_parity: true,
|
||||
http_preconditions: Some(preconditions),
|
||||
write_completion: crate::object_api::WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
};
|
||||
let mut opts = ObjectOptions::default();
|
||||
// Match migration: fixed object namespace -> durable pool metadata ->
|
||||
// actual replica namespace. The replica need not be the hash-routed set.
|
||||
let object_guard = if store.single_pool() {
|
||||
@@ -3110,9 +3105,14 @@ impl ECStore {
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let result = set
|
||||
.put_object(RUSTFS_META_BUCKET, object, &mut PutObjReader::from_vec(data), &opts)
|
||||
.await;
|
||||
let result = crate::data_movement::scanner_backlog::persist_native_scanner_pause_backlog_replica(
|
||||
set,
|
||||
data,
|
||||
preconditions,
|
||||
opts,
|
||||
"publish",
|
||||
)
|
||||
.await;
|
||||
drop(capacity_guard);
|
||||
drop(object_guard);
|
||||
result
|
||||
@@ -3747,29 +3747,37 @@ impl ECStore {
|
||||
opts: &ObjectOptions,
|
||||
no_lock: bool,
|
||||
) -> Result<usize> {
|
||||
let capacity_owner = DecommissionCapacityOwner::from_options(opts);
|
||||
match self
|
||||
.get_pool_info_existing_with_opts(bucket, object, &data_movement_pool_lookup_opts(opts, no_lock))
|
||||
.await
|
||||
{
|
||||
Ok((pinfo, _)) => Ok(pinfo.index),
|
||||
Ok((pinfo, _)) => {
|
||||
if let Some(owner) = capacity_owner {
|
||||
if self.is_decommission_capacity_target_reserved(owner, pinfo.index).await? {
|
||||
return Ok(pinfo.index);
|
||||
}
|
||||
} else {
|
||||
return Ok(pinfo.index);
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
if let Some(owner) = DecommissionCapacityOwner::from_options(opts) {
|
||||
let expected_data_bytes = opts
|
||||
.capacity_expected_data_bytes()
|
||||
.or_else(|| usize::try_from(size).ok())
|
||||
.unwrap_or_default();
|
||||
return self
|
||||
.select_decommission_capacity_target_pool(owner, expected_data_bytes)
|
||||
.await;
|
||||
}
|
||||
|
||||
self.get_available_pool_idx(bucket, object, size).await.ok_or(Error::DiskFull)
|
||||
}
|
||||
}
|
||||
if let Some(owner) = capacity_owner {
|
||||
let expected_data_bytes = opts
|
||||
.capacity_expected_data_bytes()
|
||||
.or_else(|| usize::try_from(size).ok())
|
||||
.unwrap_or_default();
|
||||
return self
|
||||
.select_decommission_capacity_target_pool(owner, expected_data_bytes)
|
||||
.await;
|
||||
}
|
||||
|
||||
self.get_available_pool_idx(bucket, object, size).await.ok_or(Error::DiskFull)
|
||||
}
|
||||
|
||||
async fn find_data_movement_target_info(
|
||||
|
||||
@@ -90,9 +90,10 @@ pub use scanner::{
|
||||
ScannerCycleScheduleStatus, ScannerPauseBacklogAlertReason, ScannerPauseBacklogPhase, ScannerPauseBacklogStatus,
|
||||
ScannerPauseBacklogThresholds, ScannerRecoveryIntentAcceptResult, ScannerRecoveryIntentConflict, ScannerRecoveryIntentRecord,
|
||||
ScannerRecoveryIntentRequest, ScannerUsageStateResetResult, accept_scanner_usage_recovery_intent,
|
||||
get_scanner_usage_recovery_intent, init_data_scanner, init_scanner_with_recovery, reset_scanner_cycle_recovery,
|
||||
reset_scanner_usage_state_for_full_rebuild, run_scanner_usage_recovery_intent, scanner_cycle_recovery_status,
|
||||
scanner_cycle_schedule_status, scanner_pause_backlog_status, scanner_recovery_actor_sha256, scanner_topology_digest,
|
||||
get_scanner_usage_recovery_intent, init_data_scanner, init_scanner_with_recovery, register_scanner_pause_backlog_retirement,
|
||||
reset_scanner_cycle_recovery, reset_scanner_usage_state_for_full_rebuild, run_scanner_usage_recovery_intent,
|
||||
scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_pause_backlog_status, scanner_recovery_actor_sha256,
|
||||
scanner_topology_digest,
|
||||
};
|
||||
pub use scanner_io::{
|
||||
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState,
|
||||
|
||||
@@ -3628,7 +3628,7 @@ pub(crate) use activity::{
|
||||
pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance};
|
||||
pub use backlog::{
|
||||
ScannerPauseBacklogAlertReason, ScannerPauseBacklogPhase, ScannerPauseBacklogStatus, ScannerPauseBacklogThresholds,
|
||||
scanner_pause_backlog_status,
|
||||
register_scanner_pause_backlog_retirement, scanner_pause_backlog_status,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test;
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -61,28 +61,43 @@ async fn setup_scanner_cycle_store_with_pool_count(
|
||||
}
|
||||
|
||||
async fn setup_scanner_cycle_store_at_path(root: &Path, seed_usage_baseline: bool, pool_count: usize) -> Arc<ECStore> {
|
||||
setup_scanner_cycle_store_at_path_with_sets(root, seed_usage_baseline, pool_count, 1).await
|
||||
}
|
||||
|
||||
pub(super) async fn setup_scanner_cycle_store_at_path_with_sets(
|
||||
root: &Path,
|
||||
seed_usage_baseline: bool,
|
||||
pool_count: usize,
|
||||
sets_per_pool: usize,
|
||||
) -> Arc<ECStore> {
|
||||
init_ecstore_config_for_scanner_tests();
|
||||
let mut pools = Vec::with_capacity(pool_count);
|
||||
for pool_index in 0..pool_count {
|
||||
let mut endpoints = Vec::new();
|
||||
for disk_index in 0..4 {
|
||||
let disk_path = root.join(format!("pool{pool_index}/disk{disk_index}"));
|
||||
tokio::fs::create_dir_all(&disk_path)
|
||||
.await
|
||||
.expect("scanner cycle test disk should be created");
|
||||
let mut endpoint =
|
||||
Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse");
|
||||
endpoint.set_pool_index(pool_index);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(disk_index);
|
||||
endpoints.push(endpoint);
|
||||
for set_index in 0..sets_per_pool {
|
||||
for disk_index in 0..4 {
|
||||
let disk_path = if sets_per_pool == 1 {
|
||||
root.join(format!("pool{pool_index}/disk{disk_index}"))
|
||||
} else {
|
||||
root.join(format!("pool{pool_index}/set{set_index}/disk{disk_index}"))
|
||||
};
|
||||
tokio::fs::create_dir_all(&disk_path)
|
||||
.await
|
||||
.expect("scanner cycle test disk should be created");
|
||||
let mut endpoint =
|
||||
Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse");
|
||||
endpoint.set_pool_index(pool_index);
|
||||
endpoint.set_set_index(set_index);
|
||||
endpoint.set_disk_index(disk_index);
|
||||
endpoints.push(endpoint);
|
||||
}
|
||||
}
|
||||
pools.push(PoolEndpoints {
|
||||
legacy: false,
|
||||
set_count: 1,
|
||||
set_count: sets_per_pool,
|
||||
drives_per_set: 4,
|
||||
endpoints: Endpoints::from(endpoints),
|
||||
cmd_line: if pool_count == 1 {
|
||||
cmd_line: if pool_count == 1 && sets_per_pool == 1 {
|
||||
"scanner-cycle-metrics".to_string()
|
||||
} else {
|
||||
format!("scanner-cycle-metrics-pool-{pool_index}")
|
||||
|
||||
@@ -28,6 +28,13 @@ pub(crate) use s3s::dto::{
|
||||
#[cfg(test)]
|
||||
pub(crate) use s3s::dto::{ExpirationStatus as EcstoreExpirationStatus, LifecycleRule as EcstoreLifecycleRule};
|
||||
|
||||
pub(crate) use rustfs_ecstore::api::data_usage::{
|
||||
MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementReplica,
|
||||
register_scanner_pause_backlog_retirement_planner,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use rustfs_ecstore::api::data_usage::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier};
|
||||
|
||||
pub(crate) use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys as EcstoreBucketTargetSys;
|
||||
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc as EcstoreLcEventSrc;
|
||||
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::{
|
||||
@@ -135,6 +142,13 @@ use rustfs_storage_api as storage_contracts;
|
||||
pub(crate) type EcstoreHealResultItem = <EcstoreStore as storage_contracts::HealOperations>::HealResultItem;
|
||||
|
||||
pub(crate) mod owner {
|
||||
pub(crate) use super::{
|
||||
MAX_SCANNER_PAUSE_BACKLOG_BYTES, ScannerPauseBacklogRetirementPlan, ScannerPauseBacklogRetirementReplica,
|
||||
register_scanner_pause_backlog_retirement_planner,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use super::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier};
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) use rustfs_ecstore::api::set_disk::test_util::hold_namespace_commit as ecstore_hold_namespace_commit;
|
||||
|
||||
|
||||
@@ -152,6 +152,7 @@ pub(crate) async fn init_startup_storage_runtime(
|
||||
readiness: Arc<GlobalReadiness>,
|
||||
instance_ctx: Arc<InstanceContext>,
|
||||
) -> Result<StartupStorageRuntime> {
|
||||
rustfs_scanner::register_scanner_pause_backlog_retirement();
|
||||
let ctx = CancellationToken::new();
|
||||
|
||||
debug!(
|
||||
@@ -195,6 +196,7 @@ pub(crate) async fn init_embedded_startup_storage_runtime(
|
||||
shutdown_token: CancellationToken,
|
||||
instance_ctx: Arc<InstanceContext>,
|
||||
) -> Result<StartupStorageRuntime> {
|
||||
rustfs_scanner::register_scanner_pause_backlog_retirement();
|
||||
let store =
|
||||
match ECStore::new_with_instance_ctx(server_addr, endpoint_pools.clone(), shutdown_token.clone(), instance_ctx).await {
|
||||
Ok(store) => store,
|
||||
|
||||
Reference in New Issue
Block a user