Compare commits

..

10 Commits

Author SHA1 Message Date
Zhengchao An 6a879be1fb test: report tier reference proof setup failures without hanging (#7608) 2026-09-10 13:35:01 +08:00
houseme 0b40d47a8a fix(scanner): bind scoped ack confirmation to generation (#7619)
Require lost scoped dirty-usage ACK reconciliation to observe a clean peer activity generation that covers the requested ACK generation. Apply the same guard in the scanner aggregation path so a stale clean activity snapshot cannot discharge pending maintenance after an uncertain ACK response.

Refs rustfs/backlog#2427

Refs rustfs/backlog#2281

Refs rustfs/backlog#2240

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 12:38:23 +08:00
houseme e292637ee0 test(scanner): collect status outcome raw evidence (#7618)
Normalize live Scanner/Heal status/outcome observations into the measured raw artifacts consumed by the G05/G06/R-D release descriptor producer.

Reject synthetic observations, incomplete required cases, and reused run/window identities before raw artifacts can enter the release bundle flow.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 12:37:11 +08:00
houseme 8b6b1e53a3 test(scanner): reject synthetic heal evidence inputs (#7617)
Harden the Scanner/Heal checkpoint restart and G14 EC evidence producers so release descriptors cannot be assembled from marked fixture, dry-run, or synthetic case inputs.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 12:20:27 +08:00
houseme 358bf6832d fix(scanner): harden heal release gates (#7614) 2026-09-10 12:03:55 +08:00
Zhengchao An 70fb8bf504 fix(ci): align release E2E gate and repair regression tests (#7607) 2026-09-10 12:03:27 +08:00
houseme 110f630a5c fix(scanner): bound native backlog fixture resources (#7609)
* fix(scanner): bound native backlog fixture resources

Keep native scanner backlog restart fixtures within low file descriptor limits by reducing the native-only disk layout, making fixture shutdown drain background work, and avoiding long-lived ECStore retention from background loops.

Convert ECStore-backed background refresh/recovery/monitor tasks to upgrade weak owners only while doing work so completed test stores release their disk graph before the next native fixture starts.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* test(tier): refresh release-merged failure fixtures

Reuse the raw legacy-transition fixture helper so delete-all expiry actually exercises Unknown transitioned history metadata after the release merge.

Rewrite restore-failure disk fixtures through per-disk xl.meta snapshots so stale destination identities can be persisted without tripping ordinary metadata update guards.

Update the SSE KMS mismatch expectation to the InvalidRequest/context_mismatch behavior now returned by the rio-v2 path.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 11:28:27 +08:00
houseme 54bd3c8e24 test(scanner): stabilize W13 release gate evidence (#7613)
* feat(nightly): publish packages as assets of the rolling 'nightly' release (sync from release) (#7593)

feat(nightly): publish packages as assets of the rolling 'nightly' release (#7592)

Replace the assets-branch scheme with a proper GitHub Release on
rustfs/auto-testing: a single 'nightly' release whose deb/rpm assets
are replaced in place on every build. This is the standard channel —
visible on the repo's Releases page, stable download URLs, no git
history growth (release assets live outside the repository).

- New scripts/release/publish_nightly_assets.sh: resolves-or-creates
  the 'nightly' release via the REST API, deletes same-name assets,
  uploads rustfs-nightly-latest.{deb,rpm}, then PATCHes the release
  body with the build provenance (ref@sha, run link, sizes, SHA256).
  Plain curl + python3, no gh CLI (the build fleet has none — #7586).
- The workflow step shrinks to invoking the script; full flow
  exercised end-to-end against the real release with probe files
  (create / upload / overwrite / download round-trip / body update).

* test(scanner): stabilize W13 release gate evidence

Treat only real raw-entry windows as replayed scanner enumeration in the restart diagnostic, so final completion rounds without raw entries are not fail-closed as raw replays.

Allow the EC8:4 multi-pool heal evidence case to select an outage object key that routes to an online pool while an entire target pool is down, preserving strict behavior for single-pool cases.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: hector <42570491+majinghe@users.noreply.github.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 10:08:55 +08:00
Zhengchao An 1643cb29f1 test: align transition fixtures and KMS context expectations (#7610) 2026-09-10 10:07:50 +08:00
cxymds bf8d5a32b6 fix(ecstore): retry contended decommission cancel target locks (#7612) 2026-09-10 09:52:43 +08:00
35 changed files with 2441 additions and 280 deletions
+1 -1
View File
@@ -1,2 +1,2 @@
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
sha256-linux=9515861be899ceb10e2e0ef93c34208bb7a7a8a7f8067a02db4cfba23270ebd6
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
@@ -22,7 +22,11 @@ mod tests {
init_logging, rustfs_binary_path,
};
use crate::storage_api::RUSTFS_META_BUCKET;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::{
error::{ProvideErrorMetadata, SdkError},
operation::put_object::PutObjectError,
primitives::ByteStream,
};
use http::Method;
use sha2::{Digest, Sha256};
use std::collections::HashSet;
@@ -586,6 +590,10 @@ mod tests {
&& operations["activeBySource"]["admin"].as_u64() == Some(1)
}
fn is_service_unavailable_put(error: &SdkError<PutObjectError>) -> bool {
error.as_service_error().and_then(ProvideErrorMetadata::code) == Some("ServiceUnavailable")
}
async fn replacement_recovery_status(
cluster: &RustFSTestClusterEnvironment,
) -> Result<serde_json::Value, Box<dyn Error + Send + Sync>> {
@@ -1342,18 +1350,41 @@ mod tests {
"replacement target must retain only its preformatted topology identity"
);
let outage_key = "cluster/written-while-node-down.bin";
let outage_payload_seed = 0xf1;
timeout(
Duration::from_secs(30),
clients[2]
.put_object()
.bucket(bucket)
.key(outage_key)
.body(ByteStream::from(deterministic_object_body(object_size_bytes, outage_payload_seed)))
.send(),
)
.await??;
let max_outage_write_attempts = if outage_target_manifest_required {
1
} else {
topology.total_drives().max(1)
};
let mut outage_key = None;
for attempt in 0..max_outage_write_attempts {
let candidate_key = if max_outage_write_attempts == 1 {
"cluster/written-while-node-down.bin".to_string()
} else {
format!("cluster/written-while-node-down-{attempt:04}.bin")
};
let put_result = timeout(
Duration::from_secs(30),
clients[2]
.put_object()
.bucket(bucket)
.key(&candidate_key)
.body(ByteStream::from(deterministic_object_body(object_size_bytes, outage_payload_seed)))
.send(),
)
.await;
match put_result {
Ok(Ok(_)) => {
outage_key = Some(candidate_key);
break;
}
Ok(Err(error)) if !outage_target_manifest_required && is_service_unavailable_put(&error) => {}
Ok(Err(error)) => return Err(error.into()),
Err(error) => return Err(error.into()),
}
}
let outage_key = outage_key
.ok_or_else(|| format!("no online pool accepted an outage object after {max_outage_write_attempts} candidates"))?;
let mut outage_peer_erasure_indices = HashSet::new();
for (node_index, node) in cluster.nodes.iter().enumerate() {
@@ -1361,7 +1392,7 @@ mod tests {
continue;
}
for (drive_index, drive) in node.data_dirs.iter().enumerate() {
let census = census_object_version_on_disk(Path::new(drive), bucket, outage_key, None)?;
let census = census_object_version_on_disk(Path::new(drive), bucket, &outage_key, None)?;
if !census.has_xl_meta {
continue;
}
@@ -1457,7 +1488,7 @@ mod tests {
"non-admin Heal is disabled, so the replacement target must remain empty before the explicit root heal"
);
assert!(
!census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?.has_xl_meta,
!census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?.has_xl_meta,
"the object written during the outage must be absent before the explicit root heal"
);
assert_eq!(
@@ -1744,10 +1775,10 @@ mod tests {
loop {
let baseline_recovered = metadata_count(&replaced_disk, bucket, &expected_manifests) == expected_manifests.len();
let outage_recovered =
!outage_target_manifest_required || object_metadata_exists_on_disk(&replaced_disk, bucket, outage_key);
!outage_target_manifest_required || object_metadata_exists_on_disk(&replaced_disk, bucket, &outage_key);
if baseline_recovered && outage_recovered {
let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?;
let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?;
let outage_census = census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?;
let pool_metadata_matches = match &expected_pool_metadata {
Some(expected) => {
census_object_version_on_disk(&replaced_disk, RUSTFS_META_BUCKET, POOL_METADATA_OBJECT, None)?
@@ -1764,7 +1795,7 @@ mod tests {
}
if Instant::now() >= heal_deadline {
let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?;
let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?;
let outage_census = census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?;
let pool_metadata =
census_object_version_on_disk(&replaced_disk, RUSTFS_META_BUCKET, POOL_METADATA_OBJECT, None)?;
let final_status = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key)
@@ -1802,7 +1833,7 @@ mod tests {
expected.key
);
}
let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?;
let outage_census = census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?;
if outage_target_manifest_required {
assert!(
outage_census.is_complete(),
@@ -1825,7 +1856,7 @@ mod tests {
.iter()
.map(|(key, _)| key.clone())
.collect::<HashSet<_>>();
assert!(expected_keys.insert(outage_key.to_string()));
assert!(expected_keys.insert(outage_key.clone()));
let node_listings = assert_all_nodes_list_exact_keys(&clients, bucket, &expected_keys).await?;
let target_client = cluster.create_s3_client(1)?;
@@ -1849,7 +1880,7 @@ mod tests {
}));
}
}
let response = target_client.get_object().bucket(bucket).key(outage_key).send().await?;
let response = target_client.get_object().bucket(bucket).key(&outage_key).send().await?;
let actual = response.body.collect().await?.into_bytes();
let expected_outage_body = deterministic_object_body(object_size_bytes, outage_payload_seed);
assert_eq!(actual.as_ref(), expected_outage_body.as_slice(), "object body changed for {outage_key}");
@@ -1860,7 +1891,7 @@ mod tests {
"expected_sha256": sha256_hex(&expected_outage_body),
"actual_sha256": sha256_hex(&actual),
"expected_physical": null,
"physical": census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?,
"physical": census_object_version_on_disk(&replaced_disk, bucket, &outage_key, None)?,
}));
}
@@ -1220,7 +1220,7 @@ impl ExpiryState {
while state.tasks_tx.len() < n {
let (tx, rx) = mpsc::channel(EXPIRY_WORKER_QUEUE_CAPACITY);
let api = api.clone();
let api = Arc::downgrade(&api);
let rx = Arc::new(tokio::sync::Mutex::new(rx));
let stats = Arc::clone(&state.stats);
let recovery_notify = Arc::clone(&state.recovery_notify);
@@ -1248,14 +1248,18 @@ impl ExpiryState {
async fn worker(
rx: &mut Receiver<Option<ExpiryOpType>>,
api: Arc<ECStore>,
api: Weak<ECStore>,
stats: Arc<ExpiryStats>,
recovery_notify: Arc<Notify>,
) {
let cancel_token = api.ctx.background_cancel_token().unwrap_or_else(|| {
let Some(initial_api) = api.upgrade() else {
return;
};
let cancel_token = initial_api.ctx.background_cancel_token().unwrap_or_else(|| {
static FALLBACK: std::sync::OnceLock<tokio_util::sync::CancellationToken> = std::sync::OnceLock::new();
FALLBACK.get_or_init(tokio_util::sync::CancellationToken::new).clone()
});
drop(initial_api);
loop {
select! {
@@ -1284,6 +1288,9 @@ impl ExpiryState {
let v = v.expect("received None after None check");
stats.decrement_pending_tasks();
let _active_task = ExpiryActiveTask::begin(Arc::clone(&stats));
let Some(api) = api.upgrade() else {
return;
};
if v.as_any().is::<ExpiryTask>() {
let v = v.as_any().downcast_ref::<ExpiryTask>().expect("ExpiryTask downcast failed");
//debug!("lifecycle expiry worker received task: {:?}", v.obj_info);
@@ -7759,8 +7766,9 @@ mod tests {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let worker_stats = Arc::clone(&stats);
let worker_notify = Arc::clone(&recovery_notify);
let worker_store = Arc::downgrade(&ecstore);
let worker = tokio::spawn(async move {
ExpiryState::worker(&mut rx, ecstore, worker_stats, worker_notify).await;
ExpiryState::worker(&mut rx, worker_store, worker_stats, worker_notify).await;
});
let oi = ObjectInfo {
bucket: "bucket".to_string(),
@@ -7861,8 +7869,9 @@ mod tests {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let worker_stats = Arc::clone(&stats);
let worker_notify = Arc::clone(&recovery_notify);
let worker_store = Arc::downgrade(&ecstore);
let worker = tokio::spawn(async move {
ExpiryState::worker(&mut rx, ecstore, worker_stats, worker_notify).await;
ExpiryState::worker(&mut rx, worker_store, worker_stats, worker_notify).await;
});
stats.increment_pending_tasks();
@@ -8031,7 +8040,7 @@ mod tests {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let worker_stats = Arc::clone(&stats);
let worker_notify = Arc::clone(&recovery_notify);
let worker_store = Arc::clone(&ecstore);
let worker_store = Arc::downgrade(&ecstore);
let worker = tokio::spawn(async move {
ExpiryState::worker(&mut rx, worker_store, worker_stats, worker_notify).await;
});
@@ -8118,7 +8127,7 @@ mod tests {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let worker_stats = Arc::clone(&stats);
let worker_notify = Arc::clone(&recovery_notify);
let worker_store = Arc::clone(&ecstore);
let worker_store = Arc::downgrade(&ecstore);
let worker = tokio::spawn(async move {
ExpiryState::worker(&mut rx, worker_store, worker_stats, worker_notify).await;
});
@@ -8225,7 +8234,7 @@ mod tests {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let worker_stats = Arc::clone(&stats);
let worker_notify = Arc::clone(&recovery_notify);
let worker_store = Arc::clone(&ecstore);
let worker_store = Arc::downgrade(&ecstore);
let worker = tokio::spawn(async move {
ExpiryState::worker(&mut rx, worker_store, worker_stats, worker_notify).await;
});
@@ -8279,8 +8288,9 @@ mod tests {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let worker_stats = Arc::clone(&stats);
let worker_notify = Arc::clone(&recovery_notify);
let worker_store = Arc::downgrade(&ecstore);
let worker = tokio::spawn(async move {
ExpiryState::worker(&mut rx, ecstore, worker_stats, worker_notify).await;
ExpiryState::worker(&mut rx, worker_store, worker_stats, worker_notify).await;
});
let oi = ObjectInfo {
bucket: format!("missing-bucket-{}", Uuid::new_v4()),
+57 -47
View File
@@ -248,10 +248,9 @@ fn validate_authoritative_object_lock_config(config: &ObjectLockConfiguration) -
}
pub async fn init_bucket_metadata_sys(api: Arc<ECStore>, buckets: Vec<String>) {
// The metadata system is inherently per-store (it holds the store handle
// and that store's bucket cache), so it lives on the store's own instance
// context (backlog#1052 S3) — a second instance initializes its own cell
// instead of panicking on the process-global one.
// The metadata system is inherently per-store, so it lives on the store's
// own instance context (backlog#1052 S3). It resolves the store through a
// weak handle so the context cache cannot keep the store and disks alive.
let instance_ctx = api.ctx.clone();
let is_dist_erasure = instance_ctx.is_dist_erasure().await;
@@ -317,18 +316,22 @@ fn start_refresh_buckets_metadata_loop(sys: Arc<RwLock<BucketMetadataSys>>) {
warn!("bucket metadata refresh loop skipped because background cancellation token is not initialized");
return;
};
let sys = Arc::downgrade(&sys);
tokio::spawn(async move {
refresh_buckets_metadata_loop(sys, cancel_token).await;
});
}
async fn refresh_buckets_metadata_loop(sys: Arc<RwLock<BucketMetadataSys>>, cancel_token: CancellationToken) {
async fn refresh_buckets_metadata_loop(sys: Weak<RwLock<BucketMetadataSys>>, cancel_token: CancellationToken) {
loop {
if !wait_refresh_interval_or_cancel(&cancel_token, BUCKET_METADATA_REFRESH_INTERVAL).await {
break;
}
refresh_buckets_metadata_once(sys.clone()).await;
let Some(sys) = sys.upgrade() else {
break;
};
refresh_buckets_metadata_once(sys).await;
}
}
@@ -455,7 +458,7 @@ pub(crate) async fn object_store_in(ctx: &crate::runtime::instance::InstanceCont
pub(crate) async fn object_store_if_initialized_in(ctx: &crate::runtime::instance::InstanceContext) -> Option<Arc<ECStore>> {
let sys = ctx.bucket_metadata_sys().or_else(get_global_bucket_metadata_sys)?;
Some(sys.read().await.api.clone())
sys.read().await.object_store_if_live()
}
pub(crate) async fn get_in(ctx: &crate::runtime::instance::InstanceContext, bucket: &str) -> Result<Arc<BucketMetadata>> {
@@ -475,7 +478,7 @@ pub(crate) async fn get_config_from_disk_with_presence_in(
bucket: &str,
) -> Result<(BucketMetadata, bool)> {
let sys = bucket_metadata_sys_of(ctx)?;
let api = sys.read().await.api.clone();
let api = sys.read().await.object_store();
load_bucket_metadata_parse_with_presence(api, bucket, true).await
}
@@ -738,7 +741,7 @@ pub async fn acquire_scanner_bucket_incarnation_fence(
) -> Result<BucketMetadataMutationGuard> {
super::utils::check_valid_bucket_name(bucket)?;
let sys = get_bucket_metadata_sys()?;
if expected_owner_id.is_nil() || sys.read().await.api.id != expected_owner_id || expected_incarnation_id.is_nil() {
if expected_owner_id.is_nil() || sys.read().await.object_store().id != expected_owner_id || expected_incarnation_id.is_nil() {
return Err(Error::other("scanner bucket incarnation owner does not match"));
}
acquire_config_write_guard_with_migration(sys, bucket, Some(expected_incarnation_id), false).await
@@ -751,7 +754,7 @@ async fn acquire_config_write_guard_with_migration(
migrate: bool,
) -> Result<BucketMetadataMutationGuard> {
let metadata_sys = sys.read().await.clone();
let lifecycle_guard = metadata_sys.api.acquire_bucket_lifecycle_read_lock(bucket).await?;
let lifecycle_guard = metadata_sys.object_store().acquire_bucket_lifecycle_read_lock(bucket).await?;
// Legacy buckets are migrated while the lifecycle fence prevents a
// same-name replacement. The second read under the write transaction is
@@ -782,7 +785,7 @@ async fn acquire_config_write_guard_with_migration(
"bucket config existence transaction validation",
async {
match metadata_sys
.api
.object_store()
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
{
@@ -802,7 +805,7 @@ async fn acquire_config_write_guard_with_migration(
Some(&transaction_guard),
bucket,
"bucket config incarnation transaction validation",
load_bucket_incarnation(metadata_sys.api.clone(), bucket),
load_bucket_incarnation(metadata_sys.object_store(), bucket),
),
)
.await?
@@ -1461,7 +1464,7 @@ pub struct BucketMetadataSys {
/// Physically missing names are TTL-bounded to limit memory under bogus
/// name floods while avoiding repeated namespace and erasure reads.
missing_buckets: moka::future::Cache<String, ()>,
api: Arc<ECStore>,
api: Weak<ECStore>,
}
impl BucketMetadataSys {
@@ -1489,12 +1492,17 @@ impl BucketMetadataSys {
.max_capacity(MISSING_BUCKET_MAX_ENTRIES)
.time_to_live(MISSING_BUCKET_TTL)
.build(),
api,
api: Arc::downgrade(&api),
}
}
pub(crate) fn object_store(&self) -> Arc<ECStore> {
self.api.clone()
self.object_store_if_live()
.expect("bucket metadata object store should still be live")
}
fn object_store_if_live(&self) -> Option<Arc<ECStore>> {
self.api.upgrade()
}
fn metadata_publish_lock(&self, bucket: &str) -> Arc<Mutex<MetadataPublishLockState>> {
@@ -1549,7 +1557,7 @@ impl BucketMetadataSys {
) -> Result<bool> {
await_bucket_namespace_operation(Some(namespace_guard), bucket, operation, async {
match self
.api
.object_store()
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
{
@@ -1566,7 +1574,7 @@ impl BucketMetadataSys {
}
async fn init_internal(&self, buckets: Vec<String>) -> Result<()> {
let count = self
.api
.object_store()
.pools
.iter()
.map(|pool| pool.disk_set.len())
@@ -1598,7 +1606,7 @@ impl BucketMetadataSys {
let mut futures = Vec::new();
for bucket in buckets.iter() {
let api = self.api.clone();
let api = self.object_store();
let bucket = bucket.clone();
futures.push(async move {
sleep(Duration::from_millis(30)).await;
@@ -1644,7 +1652,9 @@ impl BucketMetadataSys {
let bucket = bucket.clone();
futures.push(async move {
sleep(Duration::from_millis(30)).await;
let api = sys.read().await.api.clone();
let Some(api) = sys.read().await.object_store_if_live() else {
return Ok(());
};
let namespace_lock = api.new_ns_lock(&bucket, &bucket).await?;
let namespace_guard = namespace_lock
.get_read_lock(crate::set_disk::get_lock_acquire_timeout())
@@ -1685,7 +1695,7 @@ impl BucketMetadataSys {
Some(namespace_guard),
bucket,
"bucket metadata heal existence check",
self.api.bucket_exists_for_heal(bucket),
self.object_store().bucket_exists_for_heal(bucket),
)
.await?
{
@@ -1709,7 +1719,7 @@ impl BucketMetadataSys {
Some(namespace_guard),
bucket,
"bucket metadata heal",
self.api.heal_bucket(
self.object_store().heal_bucket(
bucket,
&HealOpts {
recreate: true,
@@ -1723,7 +1733,7 @@ impl BucketMetadataSys {
Some(namespace_guard),
bucket,
"bucket metadata load",
load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true),
load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true),
)
.await?;
match mode {
@@ -1903,7 +1913,7 @@ impl BucketMetadataSys {
// (backlog#1052 S7). Reading from the ambient handle instead made the
// read and the write of a single read-modify-write able to target
// different instances.
let mut bm = Box::pin(Self::load_bucket_metadata_for_update(self.api.clone(), bucket, parse)).await?;
let mut bm = Box::pin(Self::load_bucket_metadata_for_update(self.object_store(), bucket, parse)).await?;
if !bm.bucket_incarnation_sidecar || bm.bucket_incarnation_id != expected_incarnation_id {
return Err(Error::BucketNotFound(bucket.to_string()));
}
@@ -1942,7 +1952,7 @@ impl BucketMetadataSys {
where
F: FnOnce(&BucketMetadata) -> Result<Vec<u8>> + Send,
{
let mut bm = Box::pin(Self::load_bucket_metadata_for_update(self.api.clone(), bucket, true)).await?;
let mut bm = Box::pin(Self::load_bucket_metadata_for_update(self.object_store(), bucket, true)).await?;
if !bm.bucket_incarnation_sidecar || bm.bucket_incarnation_id != expected_incarnation_id {
return Err(Error::BucketNotFound(bucket.to_string()));
}
@@ -1993,7 +2003,7 @@ impl BucketMetadataSys {
/// server's metadata never leaks into the ambient (first) instance.
pub(crate) async fn persist_and_set(&self, bm: BucketMetadata) -> Result<()> {
let mut bm = bm;
bm.save_with_store(self.api.clone()).await?;
bm.save_with_store(self.object_store()).await?;
self.set(bm.name.clone(), Arc::new(bm)).await;
@@ -2001,8 +2011,8 @@ impl BucketMetadataSys {
}
async fn persist_new_and_set(&self, mut bm: BucketMetadata) -> Result<()> {
bm.save_with_store(self.api.clone()).await?;
save_bucket_incarnation(self.api.clone(), &bm.name, bm.bucket_incarnation_id).await?;
bm.save_with_store(self.object_store()).await?;
save_bucket_incarnation(self.object_store(), &bm.name, bm.bucket_incarnation_id).await?;
bm.bucket_incarnation_sidecar = true;
self.set(bm.name.clone(), Arc::new(bm)).await;
Ok(())
@@ -2019,7 +2029,7 @@ impl BucketMetadataSys {
return Err(Error::other("errInvalidArgument"));
}
load_bucket_metadata(self.api.clone(), bucket).await
load_bucket_metadata(self.object_store(), bucket).await
}
/// Reload persisted metadata under the bucket namespace generation fence.
@@ -2033,7 +2043,7 @@ impl BucketMetadataSys {
return Err(Error::other("errInvalidArgument"));
}
let namespace_lock = self.api.new_ns_lock(bucket, bucket).await?;
let namespace_lock = self.object_store().new_ns_lock(bucket, bucket).await?;
let namespace_guard = namespace_lock
.get_read_lock(crate::set_disk::get_lock_acquire_timeout())
.await?;
@@ -2056,7 +2066,7 @@ impl BucketMetadataSys {
Some(namespace_guard),
bucket,
"peer bucket metadata load",
load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true),
load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true),
)
.await?;
if !persisted {
@@ -2101,11 +2111,11 @@ impl BucketMetadataSys {
#[cfg(test)]
self.lazy_disk_loads.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let lock = self.api.new_ns_lock(bucket, bucket).await?;
let lock = self.object_store().new_ns_lock(bucket, bucket).await?;
let guard = lock.get_read_lock(crate::set_disk::get_lock_acquire_timeout()).await?;
#[cfg(test)]
if self.lazy_load_lock_probe.load(std::sync::atomic::Ordering::Relaxed) {
let competing = self.api.new_ns_lock(bucket, bucket).await?;
let competing = self.object_store().new_ns_lock(bucket, bucket).await?;
assert!(
competing.get_write_lock(Duration::from_millis(20)).await.is_err(),
"lazy metadata IO must start while the bucket namespace read lock is held"
@@ -2115,7 +2125,7 @@ impl BucketMetadataSys {
Some(&guard),
bucket,
"lazy bucket metadata load",
Box::pin(load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true)),
Box::pin(load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true)),
)
.await?;
@@ -2127,7 +2137,7 @@ impl BucketMetadataSys {
bucket,
"lazy bucket metadata existence check",
Box::pin(async {
self.api
self.object_store()
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
.map(|_| ())
@@ -2334,13 +2344,13 @@ impl BucketMetadataSys {
async fn get_bucket_incarnation_id_from_disk(&self, bucket: &str) -> Result<Uuid> {
let transaction_lock = self
.api
.object_store()
.new_ns_lock(RUSTFS_META_BUCKET, &bucket_metadata_transaction_lock_key(bucket))
.await?;
let _transaction_guard = transaction_lock
.get_read_lock(crate::set_disk::get_lock_acquire_timeout())
.await?;
let incarnation_id = load_bucket_incarnation(self.api.clone(), bucket).await?;
let incarnation_id = load_bucket_incarnation(self.object_store(), bucket).await?;
if _transaction_guard.is_lock_lost() {
return Err(Error::other(format!("bucket incarnation metadata transaction lock was lost: {bucket}")));
}
@@ -2376,7 +2386,7 @@ impl BucketMetadataSys {
async fn migrate_legacy_metadata(&self, bucket: &str) -> Result<BucketMetadataAuthority> {
let transaction_lock = self
.api
.object_store()
.new_ns_lock(RUSTFS_META_BUCKET, &bucket_metadata_transaction_lock_key(bucket))
.await?;
let _transaction_guard = transaction_lock
@@ -2401,7 +2411,7 @@ impl BucketMetadataSys {
return Err(Error::other(format!("injected Object Lock metadata disk read failure: {bucket}")));
}
let namespace_lock = self.api.new_ns_lock(bucket, bucket).await?;
let namespace_lock = self.object_store().new_ns_lock(bucket, bucket).await?;
let namespace_guard = namespace_lock
.get_read_lock(crate::set_disk::get_lock_acquire_timeout())
.await?;
@@ -2411,7 +2421,7 @@ impl BucketMetadataSys {
bucket,
"legacy bucket metadata existence check",
async {
self.api
self.object_store()
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
},
@@ -2435,7 +2445,7 @@ impl BucketMetadataSys {
Some(&namespace_guard),
bucket,
"legacy bucket metadata confirmation",
load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true),
load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true),
)
.await?;
if persisted && !metadata.bucket_incarnation_sidecar && !metadata.bucket_incarnation_id.is_nil() {
@@ -2457,20 +2467,20 @@ impl BucketMetadataSys {
}
#[cfg(test)]
if self.legacy_migration_lock_probe.load(std::sync::atomic::Ordering::Relaxed) {
let competing = self.api.new_ns_lock(bucket, bucket).await?;
let competing = self.object_store().new_ns_lock(bucket, bucket).await?;
assert!(
competing.get_write_lock(Duration::from_millis(20)).await.is_err(),
"bucket delete/recreate must not cross the legacy metadata migration fence"
);
}
save_bucket_incarnation(self.api.clone(), bucket, metadata.bucket_incarnation_id).await?;
save_bucket_incarnation(self.object_store(), bucket, metadata.bucket_incarnation_id).await?;
metadata.bucket_incarnation_sidecar = true;
if !persisted {
await_bucket_namespace_operation(
Some(&namespace_guard),
bucket,
"legacy bucket metadata migration",
metadata.save_with_store(self.api.clone()),
metadata.save_with_store(self.object_store()),
)
.await?;
}
@@ -2506,7 +2516,7 @@ impl BucketMetadataSys {
return Err(Error::other(format!("injected Object Lock metadata disk read failure: {bucket}")));
}
let namespace_lock = self.api.new_ns_lock(bucket, bucket).await?;
let namespace_lock = self.object_store().new_ns_lock(bucket, bucket).await?;
let namespace_guard = namespace_lock
.get_read_lock(crate::set_disk::get_lock_acquire_timeout())
.await?;
@@ -2515,7 +2525,7 @@ impl BucketMetadataSys {
bucket,
"bucket metadata snapshot existence check",
async {
self.api
self.object_store()
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
},
@@ -2531,7 +2541,7 @@ impl BucketMetadataSys {
Some(&namespace_guard),
bucket,
"bucket metadata authoritative snapshot",
load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true),
load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true),
)
.await?;
if persisted {
@@ -3695,7 +3705,7 @@ mod tests {
let mut stale = BucketMetadata::new("recreated-bucket");
stale.policy_config_json = b"old-generation".to_vec();
let namespace_lock = sys
.api
.object_store()
.new_ns_lock("recreated-bucket", "recreated-bucket")
.await
.expect("namespace lock should be created");
@@ -303,8 +303,16 @@ fn scanner_scoped_dirty_usage_ack_response_matches(
&& cleared_within_request
}
fn scanner_scoped_dirty_usage_ack_reconciled(activity: &ScannerPeerActivity, expected_instance_id: &str) -> bool {
activity.instance_id == expected_instance_id && activity.dirty_usage_pending == Some(false)
fn scanner_scoped_dirty_usage_ack_reconciled(
activity: &ScannerPeerActivity,
expected_instance_id: &str,
expected_generation: u64,
) -> bool {
activity.instance_id == expected_instance_id
&& activity.dirty_usage_pending == Some(false)
&& activity
.dirty_usage_generation
.is_some_and(|generation| generation >= expected_generation)
}
fn scanner_instance_id_is_valid(instance_id: &str) -> bool {
@@ -2290,6 +2298,7 @@ impl PeerRestClient {
entries: Vec<ScannerScopedDirtyUsageAckEntry>,
) -> Result<ScannerPeerActivity> {
use rustfs_protos::scoped_dirty_usage::*;
let expected_generation = entries.iter().map(|entry| entry.generation).max().unwrap_or(0);
let payloads = scanner_scoped_dirty_usage_ack_payloads(owner_id, instance_id.clone(), false, entries)?;
let ack_attempt = async {
let mut client = super::client::scanner_control_time_out_client(
@@ -2337,7 +2346,9 @@ impl PeerRestClient {
.await;
}
match self.scanner_scoped_dirty_usage_activity_confirmation().await {
Ok(activity) if scanner_scoped_dirty_usage_ack_reconciled(&activity, &instance_id) => Ok(activity),
Ok(activity) if scanner_scoped_dirty_usage_ack_reconciled(&activity, &instance_id, expected_generation) => {
Ok(activity)
}
_ => Err(err),
}
}
@@ -3176,30 +3187,48 @@ mod tests {
#[test]
fn scanner_scoped_dirty_usage_ack_reconciliation_requires_same_clean_instance() {
let activity = |instance_id: &str, pending| ScannerPeerActivity {
let activity = |instance_id: &str, generation, pending| ScannerPeerActivity {
instance_id: instance_id.to_string(),
namespace_generation: 1,
maintenance_generation: 1,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: Some([1; 32]),
data_movement_active: Some(false),
dirty_usage_generation: Some(9),
dirty_usage_generation: generation,
dirty_usage_pending: pending,
movement_generation: Some(1),
publication_blocked: Some(false),
};
assert!(scanner_scoped_dirty_usage_ack_reconciled(
&activity("0123456789abcdef0123456789abcdef", Some(false)),
"0123456789abcdef0123456789abcdef"
&activity("0123456789abcdef0123456789abcdef", Some(9), Some(false)),
"0123456789abcdef0123456789abcdef",
9
));
assert!(scanner_scoped_dirty_usage_ack_reconciled(
&activity("0123456789abcdef0123456789abcdef", Some(10), Some(false)),
"0123456789abcdef0123456789abcdef",
9
));
assert!(!scanner_scoped_dirty_usage_ack_reconciled(
&activity("0123456789abcdef0123456789abcdef", Some(true)),
"0123456789abcdef0123456789abcdef"
&activity("0123456789abcdef0123456789abcdef", Some(8), Some(false)),
"0123456789abcdef0123456789abcdef",
9
));
assert!(!scanner_scoped_dirty_usage_ack_reconciled(
&activity("fedcba9876543210fedcba9876543210", Some(false)),
"0123456789abcdef0123456789abcdef"
&activity("0123456789abcdef0123456789abcdef", None, Some(false)),
"0123456789abcdef0123456789abcdef",
9
));
assert!(!scanner_scoped_dirty_usage_ack_reconciled(
&activity("0123456789abcdef0123456789abcdef", Some(9), Some(true)),
"0123456789abcdef0123456789abcdef",
9
));
assert!(!scanner_scoped_dirty_usage_ack_reconciled(
&activity("fedcba9876543210fedcba9876543210", Some(9), Some(false)),
"0123456789abcdef0123456789abcdef",
9
));
}
+530 -20
View File
@@ -152,6 +152,7 @@ pub(crate) const DECOMMISSION_VERSION_COPY_ATTEMPTS: usize = 3;
const DECOMMISSION_COPY_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(50);
const DECOMMISSION_SOURCE_CHANGED_EXHAUSTION_LIMIT: usize = 100;
const DECOMMISSION_TERMINAL_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
const DECOMMISSION_CANCEL_TARGET_LOCK_MAX_ATTEMPTS: usize = 3;
const DECOMMISSION_DURABLE_ILM_RECEIPT_ROOT: &str = "decommission/ilm-receipts";
const DECOMMISSION_DURABLE_ILM_MANIFEST_ROOT: &str = "decommission/ilm-manifests";
const DECOMMISSION_DURABLE_ILM_RECEIPT_SCHEMA: &str = "v2";
@@ -9624,6 +9625,10 @@ struct DecommissionCapacityLockOrderBarrierState {
cancel_before_start_entered: tokio::sync::Notify,
cancel_before_start_release: tokio::sync::Notify,
cancel_before_start_paused: AtomicBool,
cancel_target_timeout_entered: tokio::sync::Notify,
cancel_target_timeout_release: tokio::sync::Notify,
cancel_target_timeout_paused: AtomicBool,
cancel_target_timeouts: AtomicUsize,
}
#[cfg(test)]
@@ -9663,6 +9668,10 @@ impl DecommissionCapacityLockOrderBarrier {
cancel_before_start_entered: tokio::sync::Notify::new(),
cancel_before_start_release: tokio::sync::Notify::new(),
cancel_before_start_paused: AtomicBool::new(false),
cancel_target_timeout_entered: tokio::sync::Notify::new(),
cancel_target_timeout_release: tokio::sync::Notify::new(),
cancel_target_timeout_paused: AtomicBool::new(false),
cancel_target_timeouts: AtomicUsize::new(0),
});
let mut slot = DECOMMISSION_CAPACITY_LOCK_ORDER_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
@@ -9784,6 +9793,21 @@ impl DecommissionCapacityLockOrderBarrier {
self.state.cancel_before_start_release.notify_one();
}
fn pause_cancel_target_timeout(&self) {
self.state.cancel_target_timeout_paused.store(true, Ordering::Release);
}
async fn wait_until_cancel_target_timeout(&self) {
tokio::time::timeout(std::time::Duration::from_secs(30), self.state.cancel_target_timeout_entered.notified())
.await
.expect("cancel should observe contention on its target capacity fence");
}
fn release_cancel_target_timeout(&self) {
self.state.cancel_target_timeout_paused.store(false, Ordering::Release);
self.state.cancel_target_timeout_release.notify_one();
}
#[cfg(feature = "test-util")]
pub(crate) fn release_owner(&self) {
self.state.owner_release.notify_one();
@@ -9826,6 +9850,7 @@ impl Drop for DecommissionCapacityLockOrderBarrier {
self.state.external_object_capacity_probe_release.notify_one();
self.state.external_object_commit_phase_release.notify_one();
self.state.cancel_before_start_release.notify_one();
self.state.cancel_target_timeout_release.notify_one();
let mut slot = DECOMMISSION_CAPACITY_LOCK_ORDER_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
@@ -9884,6 +9909,24 @@ async fn pause_decommission_cancel_before_start_gate(store_id: uuid::Uuid) {
}
}
#[cfg(test)]
async fn pause_decommission_cancel_target_timeout(store_id: uuid::Uuid) {
let barrier = DECOMMISSION_CAPACITY_LOCK_ORDER_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission capacity lock-order barrier should not be poisoned")
.as_ref()
.filter(|state| state.owner_store_id == store_id)
.cloned();
if let Some(barrier) = barrier {
barrier.cancel_target_timeouts.fetch_add(1, Ordering::AcqRel);
barrier.cancel_target_timeout_entered.notify_one();
if barrier.cancel_target_timeout_paused.load(Ordering::Acquire) {
barrier.cancel_target_timeout_release.notified().await;
}
}
}
#[cfg(test)]
fn notify_decommission_target_gate_retry(store_id: uuid::Uuid) {
let barrier = DECOMMISSION_CAPACITY_LOCK_ORDER_BARRIER
@@ -10521,6 +10564,7 @@ impl ECStore {
async fn acquire_decommission_capacity_terminal_guards(
&self,
source_pool_index: usize,
plan: Option<&DecommissionCapacityTerminalFencePlan>,
) -> Result<Vec<rustfs_lock::NamespaceLockGuard>> {
let Some(plan) = plan else {
@@ -10544,26 +10588,69 @@ impl ECStore {
"no storage pools available".to_string(),
)
})?;
let mut guards = Vec::with_capacity(plan.target_pool_indices.len());
for &target_pool_index in &plan.target_pool_indices {
let object = format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/{target_pool_index}");
let target_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, &object).await?;
let guard = target_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(|err| match err {
rustfs_lock::LockError::QuorumNotReached { required, achieved } => Error::NamespaceLockQuorumUnavailable {
mode: "write",
bucket: RUSTFS_META_BUCKET.to_string(),
object,
required,
achieved,
},
other => Error::Lock(other),
})?;
guards.push(guard);
// Retry only target acquisition, never persistence. Keep the original
// owner/cohort pinned so a remote Clear/start cannot retarget a cancel.
let mut attempt = 1;
'acquire_targets: loop {
let mut guards = Vec::with_capacity(plan.target_pool_indices.len());
for &target_pool_index in &plan.target_pool_indices {
let object = format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/{target_pool_index}");
let target_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, &object).await?;
let started = std::time::Instant::now();
match target_lock.get_write_lock(get_lock_acquire_timeout()).await {
Ok(guard) => guards.push(guard),
Err(err @ rustfs_lock::LockError::Timeout { .. }) => {
// Release the entire partial cohort before backoff or
// metadata reads; workers need these gates to settle I/O.
drop(guards);
#[cfg(test)]
pause_decommission_cancel_target_timeout(self.id).await;
if attempt >= DECOMMISSION_CANCEL_TARGET_LOCK_MAX_ATTEMPTS {
return Err(Error::Lock(err));
}
warn!(
event = EVENT_DECOMMISSION_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
state = "cancel_target_fence_retry",
pool_index = source_pool_index,
target_pool_index,
operation_id = %plan.operation_id,
generation = plan.generation,
owner_nonce = %plan.owner_nonce,
attempt,
max_attempts = DECOMMISSION_CANCEL_TARGET_LOCK_MAX_ATTEMPTS,
wait_ms = %started.elapsed().as_millis(),
error = %err,
"Decommission cancel will retry target capacity fencing"
);
tokio::time::sleep(DECOMMISSION_TERMINAL_RETRY_DELAY).await;
let save_guard = self.pool_meta_save_gate.lock().await;
let (_read_guard, snapshot) = self
.acquire_pool_meta_read_guard(&save_guard, "decommission cancel fence retry failed")
.await?;
if decommission_capacity_terminal_fence_plan(&snapshot, source_pool_index)?.as_ref() != Some(plan) {
return Err(decommission_capacity_blocked_error(
"decommission capacity owner or target cohort changed while retrying terminal fences",
));
}
attempt += 1;
continue 'acquire_targets;
}
Err(rustfs_lock::LockError::QuorumNotReached { required, achieved }) => {
return Err(Error::NamespaceLockQuorumUnavailable {
mode: "write",
bucket: RUSTFS_META_BUCKET.to_string(),
object,
required,
achieved,
});
}
Err(err) => return Err(Error::Lock(err)),
}
}
return Ok(guards);
}
Ok(guards)
}
pub(crate) async fn acquire_external_decommission_capacity_fence(
@@ -12493,7 +12580,7 @@ impl ECStore {
None
};
let _capacity_target_guards = if acquire_runtime_fence {
self.acquire_decommission_capacity_terminal_guards(terminal_fence_plan.as_ref())
self.acquire_decommission_capacity_terminal_guards(idx, terminal_fence_plan.as_ref())
.await?
} else {
Vec::new()
@@ -18736,6 +18823,429 @@ mod tests {
assert_eq!(reservation.inflight_target_physical_bytes, 0);
}
async fn start_target_fenced_cancel_test(
store: &Arc<ECStore>,
first_target_free: usize,
) -> (DecommissionCapacityTerminalFencePlan, DecommissionCanceler) {
crate::services::rebalance::promote_test_pool_meta_to_v2(store).await;
let layout = DecommissionErasureLayout { data: 1, parity: 0 };
set_decommission_capacity_info_overrides_for_test(
store.id,
vec![vec![
DecommissionPoolCapacityInfo::for_test(0, layout, 0, 10, 10),
DecommissionPoolCapacityInfo::for_test(1, layout, first_target_free, 10, 10 - first_target_free),
DecommissionPoolCapacityInfo::for_test(2, layout, 40, 40, 0),
]],
);
store
.save_current_pool_meta_for_decommission_start(&[0], Vec::new())
.await
.expect("persist an active target-fenced decommission");
let plan = decommission_capacity_terminal_fence_plan(&*store.pool_meta.read().await, 0)
.expect("active decommission should have a valid fence plan")
.expect("active decommission should retain its reservation");
assert_eq!(plan.model_version, DECOMMISSION_CAPACITY_TARGET_FENCE_MODEL_VERSION);
let canceler = DecommissionCanceler::new(CancellationToken::new());
*store.decommission_cancelers.write().await = vec![Some(canceler.clone()), None, None];
(plan, canceler)
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_cancel_waits_for_target_contention_past_one_lock_timeout() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("5"))], async {
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let (plan, canceler) = start_target_fenced_cancel_test(&store, 0).await;
assert_eq!(plan.target_pool_indices, vec![2]);
let target_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/2"))
.await
.expect("create the active migration target gate");
let target_guard = target_lock
.get_write_lock(std::time::Duration::from_secs(30))
.await
.expect("hold the target gate across the first cancel acquisition timeout");
let movement_gate = store.ctx.data_movement_operation_gate();
let movement_guard = movement_gate.read().await;
let cancel_store = Arc::clone(&store);
let mut cancel = tokio::spawn(async move { cancel_store.decommission_cancel(0).await });
tokio::time::timeout(get_lock_acquire_timeout() + std::time::Duration::from_secs(1), &mut cancel)
.await
.expect_err("one target-lock timeout must not end a legitimate cancel while its retry budget remains");
assert!(!canceler.is_cancelled(), "cancel must not signal its worker before durable fencing");
let mut durable = PoolMeta::default();
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("the active reservation must remain readable while cancel waits");
assert_eq!(
decommission_capacity_terminal_fence_plan(&durable, 0).expect("read the durable active plan"),
Some(plan),
"a failed target acquisition must not release or replace the durable reservation"
);
assert!(!durable.pools[0].decommission.as_ref().expect("active decommission").canceled);
drop(target_guard);
tokio::time::timeout(std::time::Duration::from_secs(30), canceler.token().cancelled())
.await
.expect("cancel should persist and signal its worker after target contention is released");
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("reload the committed cancellation from native replicas");
let info = durable.pools[0]
.decommission
.as_ref()
.expect("canceled state must remain durable");
assert!(info.canceled);
assert!(!info.complete && !info.failed);
let reservation = info
.capacity_reservation
.as_ref()
.expect("terminal capacity accounting must remain inspectable");
assert!(!reservation.active());
assert_eq!(reservation.pending_target_physical_bytes, 0);
assert_eq!(reservation.inflight_target_physical_bytes, 0);
tokio::time::timeout(std::time::Duration::from_millis(500), &mut cancel)
.await
.expect_err("the runtime-fenced cancel must wait for in-flight movement after signaling");
drop(movement_guard);
tokio::time::timeout(std::time::Duration::from_secs(30), cancel)
.await
.expect("cancel should return after in-flight movement quiesces")
.expect("cancel task should not panic")
.expect("the same cancel request should finish its durable terminal transition");
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_cancel_target_contention_exhausts_bounded_attempts_without_committing() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async {
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let (plan, canceler) = start_target_fenced_cancel_test(&store, 0).await;
let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id);
let target_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/2"))
.await
.expect("create the persistently contended target gate");
let target_guard = target_lock
.get_write_lock(std::time::Duration::from_secs(30))
.await
.expect("hold the target gate for every cancel attempt");
let err = tokio::time::timeout(std::time::Duration::from_secs(30), store.decommission_cancel(0))
.await
.expect("target contention must not retry indefinitely")
.expect_err("exhausting the acquisition budget must not report cancellation success");
match err {
Error::Lock(rustfs_lock::LockError::Timeout { resource, timeout }) => {
assert_eq!(resource, ".rustfs.sys/decommission/capacity-target/2@latest");
assert_eq!(timeout, std::time::Duration::from_secs(1));
}
other => panic!("expected the final typed target-lock timeout, got {other:?}"),
}
assert_eq!(barrier.state.cancel_target_timeouts.load(Ordering::Acquire), 3);
assert!(!canceler.is_cancelled());
assert!(canceler.is_active());
let mut durable = PoolMeta::default();
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("reload the reservation after the canceled request exhausted its budget");
assert_eq!(
decommission_capacity_terminal_fence_plan(&durable, 0).expect("the active plan should remain valid"),
Some(plan)
);
assert!(!durable.pools[0].decommission.as_ref().expect("active decommission").canceled);
drop(target_guard);
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_cancel_target_retry_releases_partial_cohort_and_rejects_remote_replacement() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async {
let (_temp_dirs, store, other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let (plan, _old_canceler) = start_target_fenced_cancel_test(&store, 8).await;
assert_eq!(plan.target_pool_indices, vec![1, 2]);
let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id);
barrier.pause_cancel_target_timeout();
let target_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/2"))
.await
.expect("create the second target gate");
let target_guard = target_lock
.get_write_lock(std::time::Duration::from_secs(30))
.await
.expect("block cancellation after it acquires the first target");
let cancel_store = Arc::clone(&store);
let cancel = tokio::spawn(async move { cancel_store.decommission_cancel(0).await });
barrier.wait_until_cancel_target_timeout().await;
let first_target_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/1"))
.await
.expect("create the partial-cohort release probe");
let first_target_guard = first_target_lock
.get_write_lock(std::time::Duration::from_secs(1))
.await
.expect("cancel must release its first target before retrying the blocked second target");
drop(first_target_guard);
drop(target_guard);
other_store
.reload_pool_meta()
.await
.expect("load the active generation on the remote node");
other_store
.decommission_cancel(0)
.await
.expect("the remote node should cancel the old generation");
other_store
.clear_decommission(0)
.await
.expect("the remote node should clear the old generation");
let (replacement_plan, replacement_canceler) = start_target_fenced_cancel_test(&other_store, 8).await;
assert_ne!(replacement_plan.operation_id, plan.operation_id);
barrier.release_cancel_target_timeout();
let err = tokio::time::timeout(std::time::Duration::from_secs(30), cancel)
.await
.expect("the stale request must finish without retrying a replacement operation")
.expect("the stale cancel task should not panic")
.expect_err("the original request must not cancel a remotely replaced generation");
assert!(err.to_string().contains("changed while retrying terminal fences"));
assert_eq!(barrier.state.cancel_target_timeouts.load(Ordering::Acquire), 1);
assert!(!replacement_canceler.is_cancelled());
let mut durable = PoolMeta::default();
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("reload the replacement after the stale cancel returns");
assert_eq!(
decommission_capacity_terminal_fence_plan(&durable, 0).expect("the replacement must retain its valid plan"),
Some(replacement_plan)
);
assert!(
!durable.pools[0]
.decommission
.as_ref()
.expect("replacement decommission")
.canceled
);
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_cancel_target_retry_survives_caller_abort_and_settles_inflight_mutation() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async {
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let (plan, canceler) = start_target_fenced_cancel_test(&store, 0).await;
let owner = DecommissionCapacityOwner {
source_pool_index: 0,
operation_id: plan.operation_id,
generation: plan.generation,
owner_nonce: plan.owner_nonce,
mutation_id: Some(uuid::Uuid::new_v4()),
};
let (entered_tx, entered_rx) = tokio::sync::oneshot::channel();
let (release_tx, release_rx) = tokio::sync::oneshot::channel();
let mutation_store = Arc::clone(&store);
let mutation = tokio::spawn(async move {
mutation_store
.run_decommission_capacity_admitted_mutation(2, Some(owner), Some(1), || async {
entered_tx
.send(())
.expect("the test should observe the admitted target mutation");
release_rx.await.expect("the test should release the in-flight mutation");
Ok(())
})
.await
});
tokio::time::timeout(std::time::Duration::from_secs(30), entered_rx)
.await
.expect("target mutation should reach its controlled I/O phase")
.expect("target admission should succeed before cancel starts");
let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id);
let cancel_store = Arc::clone(&store);
let cancel = tokio::spawn(async move { cancel_store.decommission_cancel(0).await });
barrier.wait_until_cancel_target_timeout().await;
assert!(!canceler.is_cancelled(), "in-flight capacity settlement must precede the cancel signal");
cancel.abort();
assert!(cancel.await.expect_err("the RPC waiter should be aborted").is_cancelled());
release_tx
.send(())
.expect("the mutation must still be alive after the caller disconnects");
tokio::time::timeout(std::time::Duration::from_secs(30), mutation)
.await
.expect("the in-flight mutation must be able to settle without a metadata lock cycle")
.expect("the mutation task should not panic")
.expect("the admitted mutation must settle before the target gate is released");
tokio::time::timeout(std::time::Duration::from_secs(30), async {
loop {
if store.pool_meta.read().await.pools[0]
.decommission
.as_ref()
.is_some_and(|info| info.canceled)
{
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
})
.await
.expect("the detached cancel transaction must finish after target contention clears");
assert!(canceler.is_cancelled());
let mut durable = PoolMeta::default();
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("reload cancellation after the RPC waiter was dropped");
let info = durable.pools[0].decommission.as_ref().expect("durable canceled decommission");
assert!(info.canceled && !info.failed && !info.complete);
let reservation = info
.capacity_reservation
.as_ref()
.expect("inspect settled capacity accounting");
assert!(!reservation.active());
assert_eq!(reservation.pending_target_physical_bytes, 0);
assert_eq!(reservation.inflight_target_physical_bytes, 0);
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_cancel_target_retry_does_not_replay_a_failed_durable_save() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async {
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let (plan, canceler) = start_target_fenced_cancel_test(&store, 0).await;
let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id);
barrier.pause_cancel_target_timeout();
let target_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/2"))
.await
.expect("create the target gate preceding the failed save");
let target_guard = target_lock
.get_write_lock(std::time::Duration::from_secs(30))
.await
.expect("force one target acquisition retry before persistence");
let save_calls = Arc::new(AtomicUsize::new(0));
let calls = Arc::clone(&save_calls);
let cancel_store = Arc::clone(&store);
let owner = canceler.clone();
let cancel = tokio::spawn(async move {
cancel_store
.decommission_cancel_transaction(0, Some(owner), true, move |_, _| async move {
calls.fetch_add(1, Ordering::AcqRel);
Err(Error::Timeout)
})
.await
});
barrier.wait_until_cancel_target_timeout().await;
assert_eq!(save_calls.load(Ordering::Acquire), 0, "persistence must wait for every target fence");
drop(target_guard);
barrier.release_cancel_target_timeout();
let err = tokio::time::timeout(std::time::Duration::from_secs(30), cancel)
.await
.expect("a persistence failure must end the cancellation transaction")
.expect("the failed-save task should not panic")
.expect_err("the injected durable-save failure must reach the caller");
assert!(matches!(err, Error::Timeout));
assert_eq!(save_calls.load(Ordering::Acquire), 1);
assert_eq!(barrier.state.cancel_target_timeouts.load(Ordering::Acquire), 1);
assert!(!canceler.is_cancelled());
assert!(canceler.is_active());
store
.ensure_pool_meta_side_effects_safe("verify ambiguous cancellation save blocks further writes")
.await
.expect_err("a failed durable save must retain the existing recovery gate");
let mut durable = PoolMeta::default();
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("the original reservation should remain readable after the injected save failure");
assert_eq!(
decommission_capacity_terminal_fence_plan(&durable, 0).expect("the old reservation should remain valid"),
Some(plan)
);
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_cancel_target_retry_converges_after_successive_contenders() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async {
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let (_plan, canceler) = start_target_fenced_cancel_test(&store, 0).await;
let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id);
barrier.pause_cancel_target_timeout();
let target_lock = store.pools[0]
.new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/2"))
.await
.expect("create the repeatedly contended target gate");
let first_guard = target_lock
.get_write_lock(std::time::Duration::from_secs(30))
.await
.expect("the first contender should own the target gate");
let cancel_store = Arc::clone(&store);
let cancel = tokio::spawn(async move { cancel_store.decommission_cancel(0).await });
barrier.wait_until_cancel_target_timeout().await;
assert_eq!(barrier.state.cancel_target_timeouts.load(Ordering::Acquire), 1);
drop(first_guard);
let second_guard = target_lock
.get_write_lock(std::time::Duration::from_secs(30))
.await
.expect("a second contender should be able to acquire between cancel attempts");
barrier.release_cancel_target_timeout();
barrier.pause_cancel_target_timeout();
barrier.wait_until_cancel_target_timeout().await;
assert_eq!(barrier.state.cancel_target_timeouts.load(Ordering::Acquire), 2);
assert!(!canceler.is_cancelled());
drop(second_guard);
barrier.release_cancel_target_timeout();
tokio::time::timeout(std::time::Duration::from_secs(30), cancel)
.await
.expect("cancel should converge after the repeated contention ends")
.expect("the cancel task should not panic")
.expect("the third target acquisition should permit a durable cancellation");
assert!(canceler.is_cancelled());
let mut durable = PoolMeta::default();
durable
.load_no_lock_from_replicas(store.pools.clone())
.await
.expect("the successful retried cancellation must survive a native metadata reload");
assert!(
durable.pools[0]
.decommission
.as_ref()
.expect("canceled decommission")
.canceled
);
assert!(
decommission_capacity_terminal_fence_plan(&durable, 0)
.expect("valid terminal metadata")
.is_none()
);
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn stale_node_cancel_cannot_replace_a_new_durable_v2_operation() {
+24 -8
View File
@@ -55,7 +55,7 @@ use rustfs_madmin::heal_commands::HealResultItem;
use rustfs_utils::{crc_hash, path::path_join_buf, sip_hash};
use std::{
collections::{HashMap, HashSet},
sync::Arc,
sync::{Arc, Weak},
};
use tokio::sync::RwLock;
use tokio::sync::broadcast::{Receiver, Sender};
@@ -308,10 +308,9 @@ impl Sets {
ctx: instance_ctx,
});
let asets = sets.clone();
let rx1 = rx.resubscribe();
tokio::spawn(async move { asets.monitor_and_connect_endpoints(rx1).await });
let weak_sets = Arc::downgrade(&sets);
tokio::spawn(async move { Self::monitor_and_connect_endpoints_task(weak_sets, rx1).await });
Ok(sets)
}
@@ -326,12 +325,26 @@ impl Sets {
&self.ctx
}
pub async fn monitor_and_connect_endpoints(&self, mut rx: Receiver<()>) {
tokio::time::sleep(Duration::from_secs(5)).await;
async fn monitor_and_connect_endpoints_task(sets: Weak<Sets>, mut rx: Receiver<()>) {
let startup_delay = tokio::time::sleep(Duration::from_secs(5));
tokio::pin!(startup_delay);
tokio::select! {
_ = &mut startup_delay => {}
_ = rx.recv() => {
warn!("monitor_and_connect_endpoints ctx cancelled");
return;
}
}
info!("start monitor_and_connect_endpoints");
self.connect_disks().await;
let Some(current) = sets.upgrade() else {
warn!("monitor_and_connect_endpoints exit");
return;
};
current.connect_disks().await;
drop(current);
// TODO(backlog): make monitor_and_connect interval configurable instead of hardcoded 15s
let mut interval = tokio::time::interval(Duration::from_secs(15));
@@ -339,7 +352,10 @@ impl Sets {
tokio::select! {
_= interval.tick()=>{
// debug!("tick...");
self.connect_disks().await;
let Some(current) = sets.upgrade() else {
break;
};
current.connect_disks().await;
interval.reset();
},
+106 -4
View File
@@ -6433,9 +6433,48 @@ impl TierConfigMgr {
}
pub(crate) async fn refresh_tier_config_handle(handle: Arc<RwLock<Self>>, api: Arc<ECStore>) {
Self::refresh_tier_config_handle_with(handle, api).await;
Self::refresh_tier_config_handle_with_weak(handle, Arc::downgrade(&api)).await;
}
async fn refresh_tier_config_handle_with_weak(handle: Arc<RwLock<Self>>, api: Weak<ECStore>) {
// The periodic refresh remains the recovery fallback; committed mutations
// notify this worker so a successful peer commit converges immediately.
let mutation_refresh = Self::mutation_refresh_notifier(&handle).await;
let r = rand::rng().random_range(0.0..1.0);
let rand_interval = || Duration::from_secs((r * 60_f64).round() as u64);
let refresh_interval = TIER_CFG_REFRESH + rand_interval();
let mut t = delayed_tier_refresh_interval(refresh_interval);
loop {
select! {
_ = t.tick() => {
let Some(api) = Weak::upgrade(&api) else {
return;
};
if let Err(err) = Self::reload_handle_with(&handle, api).await {
warn!(
event = EVENT_TIER_CONFIG_REFRESH,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_TIER,
trigger = "periodic",
result = "failed",
error = ?err,
"tier configuration refresh"
);
}
}
_ = mutation_refresh.notified() => {
let Some(api) = Weak::upgrade(&api) else {
return;
};
Self::reload_after_committed_mutation(&handle, api).await;
}
}
t.reset();
}
}
#[allow(dead_code, reason = "used by focused tier refresh tests and non-ECStore generic harnesses")]
pub(crate) async fn refresh_tier_config_handle_with<S>(handle: Arc<RwLock<Self>>, api: Arc<S>)
where
S: EcstoreObjectIO
@@ -17540,6 +17579,66 @@ mod tests {
assert!(current.tiers.contains_key("COLD-B"));
}
async fn wait_for_reference_proof_barrier(
barrier: &TierDriverBuildBarrier,
update: &mut tokio::task::JoinHandle<std::result::Result<(), TierConfigUpdateError>>,
) -> std::result::Result<(), String> {
tokio::select! {
biased;
result = &mut *update => Err(format!("tier update exited before the reference proof barrier: {result:?}")),
() = barrier.arrived.notified() => Ok(()),
() = tokio::time::sleep(Duration::from_secs(30)) => {
// Aborting the caller does not stop its owned mutation task.
// Let a late arrival pass the test-only barrier.
barrier.release.add_permits(1);
update.abort();
Err("timed out waiting for the reference proof barrier".to_string())
}
}
}
#[tokio::test]
#[serial_test::serial]
async fn reference_proof_barrier_reports_update_failure_before_arrival() {
let manager = TierConfigMgr::new();
let store = Arc::new(CasConfigStore::default());
let mut persisted = empty_mgr();
persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A"));
persisted
.save_tiering_config_if_current(store.clone(), None)
.await
.expect("early update failure fixture should persist");
let barrier = tier_reference_proof_test_barrier();
let scoped_barrier = barrier.clone();
let factory: TierDriverTestFactory =
Arc::new(|_| Err(AdminError::msg("injected driver initialization failure before reference proof")));
let mut update = tokio::spawn(async move {
TIER_REFERENCE_PROOF_TEST_BARRIER
.scope(
scoped_barrier,
TIER_DRIVER_TEST_FACTORY.scope(
factory,
TIER_MUTATION_TEST_PEERS.scope(
Vec::new(),
TierConfigMgr::update_candidate_with_config_lock(
&manager,
store,
TierCandidateMutation::Remove("COLD-A".to_string(), true),
),
),
),
)
.await
});
let err = tokio::time::timeout(Duration::from_secs(5), wait_for_reference_proof_barrier(&barrier, &mut update))
.await
.expect("an early update failure should be observed without waiting for the barrier deadline")
.expect_err("a failed update cannot reach the reference proof barrier");
assert!(err.contains("Mutation"), "{err}");
assert!(err.contains("injected driver initialization failure before reference proof"), "{err}");
}
#[tokio::test]
#[serial_test::serial]
async fn reference_proof_rejects_a_changed_prepared_fence_revision_before_publish() {
@@ -17560,7 +17659,7 @@ mod tests {
let scoped_barrier = barrier.clone();
let update_manager = manager.clone();
let update_store = store.clone();
let update = tokio::spawn(async move {
let mut update = tokio::spawn(async move {
TIER_REFERENCE_PROOF_TEST_BARRIER
.scope(
scoped_barrier,
@@ -17575,7 +17674,9 @@ mod tests {
)
.await
});
barrier.arrived.notified().await;
wait_for_reference_proof_barrier(&barrier, &mut update)
.await
.expect("tier update should reach the reference proof barrier");
let unrelated = prepared_remove_intent("COLD-B", uuid::Uuid::from_u128(0x2237));
TierConfigMgr::apply_prepared_mutation_intent_block(&manager, &unrelated)
@@ -17583,8 +17684,9 @@ mod tests {
.expect("an unrelated prepared fence should advance the runtime revision");
barrier.release.add_permits(1);
let err = update
let err = tokio::time::timeout(Duration::from_secs(30), update)
.await
.expect("tier update should finish after the reference proof barrier releases")
.expect("tier update task should join")
.expect_err("a reference proof cannot authorize publication across a fence revision change");
let TierConfigUpdateError::Publish(err) = err else {
+97 -8
View File
@@ -13807,7 +13807,17 @@ mod transition_commit_failure_tests {
#[tokio::test]
#[serial_test::serial]
async fn restore_failure_after_snapshot_cleans_exact_generation_and_returns_primary_error() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
assert_restore_failure_cleanup_boundary(true).await;
}
#[tokio::test]
#[serial_test::serial]
async fn restore_failure_after_snapshot_preserves_corrupt_known_transition_metadata() {
assert_restore_failure_cleanup_boundary(false).await;
}
async fn assert_restore_failure_cleanup_boundary(legacy_unknown: bool) {
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "restore-post-snapshot-cleanup-bucket";
let object = "object.bin";
for disk in &disk_stores {
@@ -13816,7 +13826,15 @@ mod transition_commit_failure_tests {
let mut reader = PutObjReader::from_vec(b"post-snapshot cleanup source".repeat(1024));
let original = set_disks
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
write_completion: WriteCompletion::TailDrained,
..Default::default()
},
)
.await
.expect("source object should be written");
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
@@ -13857,16 +13875,70 @@ mod transition_commit_failure_tests {
.await
.expect("transitioned metadata should be readable")
.into_owned();
let known_state = source_fi.transition_version_state;
assert_ne!(known_state, rustfs_filemeta::TransitionVersionState::Unknown);
source_fi.metadata.extend(restore_metadata(operation_id, true));
rustfs_utils::http::insert_str(
&mut source_fi.metadata,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"invalid".to_string(),
);
set_disks
.update_object_meta(bucket, object, source_fi, &online_disks)
.await
.expect("invalid backend identity fixture should be persisted");
.expect("restore markers should be persisted");
// Normal writes reject damage to a reconciled binding. Model on-disk
// corruption directly, with and without the legacy missing-state field.
let mut corrupted_metadata = Vec::new();
for temp_dir in &temp_dirs {
let metadata_path = temp_dir.path().join(bucket).join(object).join(STORAGE_FORMAT_FILE);
let encoded = tokio::fs::read(&metadata_path)
.await
.expect("transition metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("transition metadata should decode");
let (version_index, mut version) = metadata
.find_version(original.version_id)
.expect("transitioned version should exist");
let object_meta = version.object.as_mut().expect("transitioned version should be an object");
rustfs_utils::http::insert_bytes(
&mut object_meta.meta_sys,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
b"invalid".to_vec(),
);
if legacy_unknown {
rustfs_utils::http::remove_bytes(
&mut object_meta.meta_sys,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
);
}
metadata.versions[version_index] =
rustfs_filemeta::FileMetaShallowVersion::try_from(version).expect("corrupt fixture should re-encode");
tokio::fs::write(&metadata_path, metadata.marshal_msg().expect("corrupt fixture should encode"))
.await
.expect("corrupt fixture should be written");
let persisted = tokio::fs::read(&metadata_path)
.await
.expect("corrupt fixture should be readable");
let fixture = FileMeta::load(&persisted)
.expect("corrupt fixture should decode")
.find_version(original.version_id)
.expect("corrupt version should exist")
.1
.into_fileinfo(bucket, object, true)
.expect("corrupt version should decode");
assert_eq!(
fixture.transition_version_state,
if legacy_unknown {
rustfs_filemeta::TransitionVersionState::Unknown
} else {
known_state
}
);
assert_eq!(
rustfs_utils::http::get_str(&fixture.metadata, rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID),
Some("invalid".to_string())
);
for (key, value) in restore_metadata(operation_id, true) {
assert_eq!(fixture.metadata.get(&key), Some(&value), "fixture must retain restore marker {key}");
}
corrupted_metadata.push((metadata_path, persisted));
}
set_disks.invalidate_get_object_metadata_cache(bucket, object).await;
let mut opts = ObjectOptions::default();
@@ -13889,6 +13961,23 @@ mod transition_commit_failure_tests {
.await
.expect("cleanup should leave the transitioned object readable");
assert_eq!(cleaned.transitioned_object.status, TRANSITION_COMPLETE);
if !legacy_unknown {
// Known bindings with corrupt identities must be repaired before
// cleanup; rejection must preserve both the binding and markers.
for (key, value) in restore_metadata(operation_id, true) {
assert_eq!(cleaned.user_defined.get(&key), Some(&value), "cleanup must preserve restore marker {key}");
}
for (metadata_path, before) in corrupted_metadata {
assert_eq!(
tokio::fs::read(metadata_path)
.await
.expect("rejected cleanup metadata should remain readable"),
before,
"rejected cleanup must leave corrupt known metadata unchanged"
);
}
return;
}
assert!(!cleaned.user_defined.contains_key(s3s::header::X_AMZ_RESTORE.as_str()));
assert!(
rustfs_utils::http::get_str(cleaned.user_defined.as_ref(), rustfs_utils::http::SUFFIX_RESTORE_OPERATION_ID,)
+37 -31
View File
@@ -15661,25 +15661,24 @@ mod tests {
assert!(deleted[0].found, "the aggregate error must retain the committed pool result");
drop(injection);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let mut metadata_absent = true;
for pool in &store.pools {
metadata_absent &= pool
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("aggregate-error cleanup metadata should remain readable")
.is_none();
}
if metadata_absent && backend.remove_count().await == 1 {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("aggregate failure must not suppress committed receipt dispatch");
// Exact reads can see subquorum metadata while workers remove each
// disk's free version. Inspect the final state after cleanup drains.
wait_for_expiry_workers_idle(&store).await;
for pool in &store.pools {
assert!(
pool.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("aggregate-error cleanup metadata should remain readable")
.is_none(),
"aggregate failure must not suppress committed receipt cleanup"
);
}
assert_eq!(
backend.remove_count().await,
1,
"committed receipts must remove the shared remote object once"
);
assert_eq!(backend.object_count().await, 0, "the shared remote object should be removed exactly once");
store
.delete_bucket(bucket, &DeleteBucketOptions::default())
@@ -19057,27 +19056,34 @@ mod tests {
.await
.expect("transition metadata should be readable");
let mut metadata = FileMeta::load(&encoded).expect("transition metadata should decode");
let mut transitioned = metadata
.get_all_file_info_versions(bucket, object, true)
.expect("transitioned versions should decode")
.versions
.into_iter()
.find(|version| version.version_id == history.version_id)
let (version_index, mut transitioned) = metadata
.find_version(history.version_id)
.expect("transitioned history should exist");
transitioned.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown;
rustfs_utils::http::metadata_compat::remove_str(
&mut transitioned.metadata,
// Rewrite the serialized record to model legacy metadata;
// ordinary writes preserve an already reconciled state.
rustfs_utils::http::metadata_compat::remove_bytes(
&mut transitioned.object.as_mut().expect("history should be an object").meta_sys,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE,
);
metadata
.add_version(transitioned)
.expect("unknown state should replace the transitioned version");
metadata.versions[version_index] = rustfs_filemeta::FileMetaShallowVersion::try_from(transitioned)
.expect("legacy history should re-encode");
tokio::fs::write(
&metadata_path,
metadata.marshal_msg().expect("unknown transition metadata should encode"),
)
.await
.expect("unknown transition metadata should be written");
let encoded = tokio::fs::read(&metadata_path)
.await
.expect("legacy transition metadata should be readable");
let legacy = FileMeta::load(&encoded)
.expect("legacy transition metadata should decode")
.find_version(history.version_id)
.expect("legacy history should exist")
.1
.into_fileinfo(bucket, object, true)
.expect("legacy history should decode");
assert_eq!(legacy.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown);
}
let lifecycle_event = crate::bucket::lifecycle::lifecycle::Event {
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
@@ -1321,6 +1321,38 @@ mod tests {
.expect("token auth must map to a source");
}
/// A pending acquisition cannot take the returned-error fallback: the
/// outer login policy must cut it off before publishing a client.
#[tokio::test(start_paused = true)]
async fn test_stalled_initial_login_is_bounded_by_the_attempt_timeout() {
let state = Arc::new(ScriptedState::default());
let source = ScriptedSource {
state: state.clone(),
ttl: Duration::ZERO,
renewable: false,
login_delay: Duration::from_secs(60),
};
let policy = test_policy(Duration::from_secs(10), Duration::from_secs(5));
let attempt_timeout = policy.retry.attempt_timeout;
let started = Instant::now();
let result = VaultCredentialProvider::new(test_settings(), Box::new(source), policy).await;
assert!(
matches!(&result, Err(KmsError::OperationTimedOut { message }) if message.starts_with("vault_login attempt 1 timed out")),
"a stalled login must return its typed timeout without publishing a client"
);
assert_eq!(
started.elapsed(),
attempt_timeout,
"login must consume exactly one virtual attempt budget"
);
assert_eq!(state.login_calls.load(Ordering::SeqCst), 0, "the acquisition must not complete");
assert_eq!(
state.renew_calls.load(Ordering::SeqCst),
0,
"failed initialization must not start renewal"
);
}
/// backlog#2369 P3: `vault token create` defaults to a 768-hour TTL, so
/// hard-coding "no lease" for token auth left the renewal task unstarted
/// and turned a healthy cluster into one that answers 403 a month later.
+200 -27
View File
@@ -14,8 +14,7 @@
//! Fault-injection matrix for the Vault backend operation policy.
//!
//! Offline cases run against locally injected transport faults (a listener
//! that never responds) — deterministic, no external
//! Offline cases run against locally injected HTTP and transport faults — no external
//! dependencies. Real-Vault cases are `#[ignore]`d and need a dev Vault
//! (default `http://127.0.0.1:8200`, override with `RUSTFS_KMS_VAULT_ADDR`).
//!
@@ -38,9 +37,100 @@ use rustfs_kms::backends::vault::VaultKmsBackend;
use rustfs_kms::{
BackendConfig, DescribeKeyRequest, KmsBackend as KmsBackendKind, KmsConfig, KmsError, VaultAuthMethod, VaultConfig,
};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpListener;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
const OPERATIONS_TOTAL: &str = "rustfs_kms_backend_operations_total";
const ATTEMPT_FAILURES_TOTAL: &str = "rustfs_kms_backend_attempt_failures_total";
const LOGIN: &str = "vault_login";
const READ_KEY: &str = "vault_kv2_read_key";
const LOOKUP_REQUEST: &str = "GET /v1/auth/token/lookup-self HTTP/1.1";
/// Unlike the unit-test scripted Vault, this fixture records the credential
/// probe too and can fail it independently of the subsequent key request.
/// `None` parks a connection without responding; extra requests receive 599.
struct FaultVault {
address: String,
requests: mpsc::UnboundedReceiver<String>,
task: JoinHandle<()>,
}
impl FaultVault {
async fn serve(responses: Vec<Option<(u16, serde_json::Value)>>) -> Self {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fault-injection Vault");
let address = format!("http://{}", listener.local_addr().expect("fault-injection Vault address"));
let (recorded, requests) = mpsc::unbounded_channel();
let task = tokio::spawn(async move {
let mut responses = responses.into_iter();
let mut parked = Vec::new();
loop {
let (stream, _) = listener.accept().await.expect("accept Vault request");
let mut stream = BufReader::new(stream);
let mut line = String::new();
assert_ne!(stream.read_line(&mut line).await.expect("read request line"), 0);
recorded.send(line.trim_end().to_string()).expect("record Vault request");
loop {
line.clear();
assert_ne!(stream.read_line(&mut line).await.expect("read request header"), 0);
if line == "\r\n" {
break;
}
}
let mut stream = stream.into_inner();
let response = responses
.next()
.unwrap_or_else(|| Some((599, serde_json::json!({"errors": ["unexpected Vault request"]}))));
if let Some((status, body)) = response {
let body = body.to_string();
let response = format!(
"HTTP/1.1 {status} Scripted\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len()
);
stream.write_all(response.as_bytes()).await.expect("write Vault response");
stream.shutdown().await.expect("close Vault response");
} else {
parked.push(stream);
}
}
});
Self { address, requests, task }
}
async fn finish(&mut self) {
assert!(self.requests.try_recv().is_err(), "no unexpected requests may remain");
self.task.abort();
let error = (&mut self.task).await.expect_err("fault server runs until aborted");
assert!(error.is_cancelled(), "fault server must not panic: {error}");
}
}
impl Drop for FaultVault {
fn drop(&mut self) {
self.task.abort();
}
}
fn healthy_token_lookup() -> serde_json::Value {
serde_json::json!({
"data": {
"accessor": "fault-injection-accessor",
"creation_time": 1_700_000_000u64,
"creation_ttl": 0,
"display_name": "token",
"entity_id": "",
"explicit_max_ttl": 0,
"id": "unused",
"num_uses": 0,
"orphan": true,
"path": "auth/token/create",
"policies": ["default"],
"renewable": false,
"ttl": 0
}
})
}
fn vault_config(address: &str, token: &str) -> VaultConfig {
VaultConfig {
@@ -125,39 +215,112 @@ fn counter_value(snapshot: &[MetricEntry], name: &str, labels: &[(&str, &str)])
fn stalled_connection_is_cut_off_by_the_attempt_timeout() {
let snapshot = record_metrics(|| {
Box::pin(async move {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
let mut vault = FaultVault::serve(vec![Some((200, healthy_token_lookup())), None]).await;
let attempt_timeout = Duration::from_millis(250);
let client = VaultKmsBackend::new(kms_config(vault_config(&vault.address, "unused"), attempt_timeout, 1))
.await
.expect("bind stall listener");
let address = format!("http://{}", listener.local_addr().expect("stall listener addr"));
// Accept and park every connection without ever responding.
tokio::spawn(async move {
let mut parked = Vec::new();
loop {
let Ok((socket, _)) = listener.accept().await else { return };
parked.push(socket);
}
});
.expect("the token lookup must succeed before injecting the stalled key read");
assert_eq!(vault.requests.try_recv().as_deref(), Ok(LOOKUP_REQUEST));
let client = VaultKmsBackend::new(kms_config(vault_config(&address, "unused"), Duration::from_millis(250), 1))
.await
.expect("client construction performs no network calls");
let error = KmsBackendTrait::describe_key(&client, describe_key_request("fault-injection-stalled"))
let read = KmsBackendTrait::describe_key(&client, describe_key_request("fault-injection-stalled"));
tokio::pin!(read);
tokio::select! {
request = vault.requests.recv() => assert_eq!(
request.as_deref(),
Some("GET /v1/secret/data/rustfs/kms/fault-injection/fault-injection-stalled? HTTP/1.1")
),
result = &mut read => panic!("the key request must reach the stall listener: {result:?}"),
}
// Pause only after real loopback I/O reaches the intended request;
// otherwise auto-advancing time could expire the login instead.
tokio::time::pause();
let stalled_at = tokio::time::Instant::now();
// Tokio rounds timer deadlines up to the next millisecond.
let virtual_step = attempt_timeout + Duration::from_millis(1);
tokio::time::advance(virtual_step).await;
let error = tokio::time::timeout(Duration::from_secs(1), read)
.await
.expect("the attempt timer must resolve without further network activity")
.expect_err("a stalled request must be cut off by the attempt timeout");
assert_eq!(
stalled_at.elapsed(),
virtual_step,
"the read must resolve within the attempt budget plus one timer tick"
);
assert!(
matches!(error, KmsError::OperationTimedOut { .. } | KmsError::BackendError { .. }),
"got {error:?}"
);
vault.finish().await;
})
});
// The policy timer reports attempt_timeout; the client-level HTTP timeout
// surfaces as a connection-class failure. Either way it is exactly one
// attempt that was cut off.
let cut_off = counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "attempt_timeout")])
+ counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "retryable_conn")]);
let cut_off = counter_value(
&snapshot,
ATTEMPT_FAILURES_TOTAL,
&[("operation", READ_KEY), ("error_class", "attempt_timeout")],
) + counter_value(
&snapshot,
ATTEMPT_FAILURES_TOTAL,
&[("operation", READ_KEY), ("error_class", "retryable_conn")],
);
assert_eq!(cut_off, 1, "the single budgeted attempt must be cut off by a timeout");
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]), 1);
assert_eq!(counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("operation", LOGIN)]), 0);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("operation", LOGIN), ("outcome", "success")]),
1
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("operation", READ_KEY), ("outcome", "budget_exhausted")]),
1
);
}
/// A returned lookup error degrades lease discovery, not Vault authorization:
/// the subsequent forbidden key read must still fail once, without retrying.
#[test]
fn token_lookup_errors_do_not_bypass_key_authorization() {
for lookup_status in [403, 503] {
let snapshot = record_metrics(|| {
Box::pin(async move {
let mut vault = FaultVault::serve(vec![
Some((lookup_status, serde_json::json!({"errors": ["token lookup unavailable"]}))),
Some((403, serde_json::json!({"errors": ["permission denied"]}))),
])
.await;
let client = VaultKmsBackend::new(kms_config(vault_config(&vault.address, "unused"), Duration::from_secs(5), 3))
.await
.expect("a returned token lookup error must preserve static-token fallback");
assert_eq!(vault.requests.try_recv().as_deref(), Ok(LOOKUP_REQUEST));
let error = KmsBackendTrait::describe_key(&client, describe_key_request("fault-injection-forbidden"))
.await
.expect_err("lease discovery fallback must not authorize a forbidden key read");
assert!(matches!(error, KmsError::BackendError { .. }), "got {error:?}");
assert_eq!(
vault.requests.try_recv().as_deref(),
Ok("GET /v1/secret/data/rustfs/kms/fault-injection/fault-injection-forbidden? HTTP/1.1")
);
vault.finish().await;
})
});
assert_eq!(counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("operation", LOGIN)]), 0);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("operation", LOGIN), ("outcome", "success")]),
1
);
assert_eq!(counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("operation", READ_KEY)]), 1);
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("operation", READ_KEY), ("error_class", "fatal")]),
1
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("operation", READ_KEY), ("outcome", "fatal")]),
1
);
}
}
fn real_vault_address() -> String {
@@ -174,7 +337,7 @@ fn real_vault_invalid_token_is_fatal_and_never_retried() {
let config = vault_config(&real_vault_address(), "fault-injection-invalid-token");
let client = VaultKmsBackend::new(kms_config(config, Duration::from_secs(5), 3))
.await
.expect("client construction performs no network calls");
.expect("a returned token lookup error must preserve static-token fallback");
let error = KmsBackendTrait::describe_key(&client, describe_key_request("fault-injection-forbidden"))
.await
.expect_err("an invalid token must be rejected");
@@ -183,13 +346,20 @@ fn real_vault_invalid_token_is_fatal_and_never_retried() {
});
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "fatal")]),
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("operation", READ_KEY), ("error_class", "fatal")]),
1,
"a 403 must be observed by exactly one attempt"
);
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "fatal")]), 1);
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "retryable_status")]),
counter_value(&snapshot, OPERATIONS_TOTAL, &[("operation", READ_KEY), ("outcome", "fatal")]),
1
);
assert_eq!(
counter_value(
&snapshot,
ATTEMPT_FAILURES_TOTAL,
&[("operation", READ_KEY), ("error_class", "retryable_status")]
),
0,
"an auth failure must never be classified as retryable"
);
@@ -207,7 +377,7 @@ fn real_vault_missing_key_is_resolved_in_one_attempt() {
let config = vault_config(&real_vault_address(), &token);
let client = VaultKmsBackend::new(kms_config(config, Duration::from_secs(5), 3))
.await
.expect("client construction performs no network calls");
.expect("static-token initialization must complete against the running Vault");
let error = KmsBackendTrait::describe_key(&client, describe_key_request("fault-injection-definitely-missing"))
.await
.expect_err("a missing key must resolve to key-not-found");
@@ -216,9 +386,12 @@ fn real_vault_missing_key_is_resolved_in_one_attempt() {
});
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "fatal")]),
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("operation", READ_KEY), ("error_class", "fatal")]),
1,
"a 404 must be observed by exactly one attempt"
);
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "fatal")]), 1);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("operation", READ_KEY), ("outcome", "fatal")]),
1
);
}
+3 -3
View File
@@ -96,14 +96,14 @@ pub use scanner::{
scanner_topology_digest,
};
pub use scanner_io::{
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageClearObserver, ScannerDirtyUsageSnapshot,
ScannerDirtyUsageState, ScannerDurableDirtyUsageReplayEntry, ScannerDurableDirtyUsageReplayError,
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageClearObserver, ScannerDirtyUsageMutationObserver,
ScannerDirtyUsageSnapshot, ScannerDirtyUsageState, ScannerDurableDirtyUsageReplayEntry, ScannerDurableDirtyUsageReplayError,
ScannerDurableDirtyUsageReplayRecord, ScannerDurableDirtyUsageReplayScope, acknowledge_dirty_usage_generation,
acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, encode_durable_dirty_usage_producer_replay_record,
record_dirty_usage_bucket, record_dirty_usage_bucket_from_producer, record_dirty_usage_bucket_from_producers,
record_dirty_usage_object, record_dirty_usage_object_from_producer, record_scanner_maintenance_change,
replay_durable_dirty_usage_producer_record, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state,
scanner_maintenance_generation, set_scanner_dirty_usage_clear_observer,
scanner_maintenance_generation, set_scanner_dirty_usage_clear_observer, set_scanner_dirty_usage_mutation_observer,
};
pub use segment_invalidation::SegmentInvalidationProducerIdentity;
pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER};
+17 -2
View File
@@ -105,8 +105,14 @@ pub(super) fn remote_dirty_usage_acknowledgement_loss_reconciled(
if !acknowledged_hosts.insert(acknowledgement.host.as_str()) {
return false;
}
scanner_activity_dirty_usage_state_for_host(&activity_after_error, &acknowledgement.host)
.is_some_and(|(instance_id, _generation, pending)| instance_id == acknowledgement.instance_id && !pending)
let Some(expected_generation) = acknowledgement.expected_dirty_usage_generation() else {
return false;
};
scanner_activity_dirty_usage_state_for_host(&activity_after_error, &acknowledgement.host).is_some_and(
|(instance_id, generation, pending)| {
instance_id == acknowledgement.instance_id && generation >= expected_generation && !pending
},
)
})
}
@@ -509,6 +515,15 @@ pub(crate) enum ScannerDirtyUsageAcknowledgementKind {
},
}
impl ScannerDirtyUsageAcknowledgement {
fn expected_dirty_usage_generation(&self) -> Option<u64> {
match &self.kind {
ScannerDirtyUsageAcknowledgementKind::Generation(generation) => Some(*generation),
ScannerDirtyUsageAcknowledgementKind::Scoped { entries, .. } => entries.iter().map(|entry| entry.generation).max(),
}
}
}
impl From<ScannerDirtyUsageAcknowledgement> for crate::storage_api::EcstoreScannerDirtyUsageAcknowledgement {
fn from(acknowledgement: ScannerDirtyUsageAcknowledgement) -> Self {
match acknowledgement.kind {
+97 -22
View File
@@ -1798,6 +1798,8 @@ pub(super) fn scanner_pause_backlog_now() -> u64 {
mod tests {
use super::*;
const NATIVE_RETIREMENT_DRIVES_PER_SET: usize = 2;
fn run_native_retirement_test<C, F>(case: C)
where
C: FnOnce() -> F + Send + 'static,
@@ -1825,10 +1827,29 @@ mod tests {
async fn native_retirement_store() -> (tempfile::TempDir, Arc<ECStore>) {
register_scanner_pause_backlog_retirement();
let root = tempfile::tempdir().expect("native retirement fixture directory");
let store = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await;
let store = super::super::tests::setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root.path(),
false,
3,
2,
NATIVE_RETIREMENT_DRIVES_PER_SET,
false,
)
.await;
(root, store)
}
async fn shutdown_native_retirement_store(store: Arc<ECStore>) {
if let Some(token) = store.background_cancel_token() {
token.cancel();
}
drop(store);
for _ in 0..8 {
tokio::task::yield_now().await;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
async fn native_replica_bytes(set: &SetDisks) -> (Vec<u8>, String) {
let mut reader = set
.get_object_reader(
@@ -1871,7 +1892,15 @@ mod tests {
register_scanner_pause_backlog_retirement();
let root = tempfile::tempdir().unwrap();
let old = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, old_pool_count, 2).await;
let old = super::super::tests::setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root.path(),
false,
old_pool_count,
2,
NATIVE_RETIREMENT_DRIVES_PER_SET,
false,
)
.await;
let fault = unstable_source
.then(|| NativeScannerPauseBacklogWriteFault::fail_before_write(Arc::clone(&old.pools[0].disk_set[0]), "publish", 2));
let now = unix_now();
@@ -1889,10 +1918,14 @@ mod tests {
} else {
controller.observe(observation(now + 1, true, 4)).await;
controller.observe(observation(now + 2, false, 4)).await;
assert!(matches!(
controller.begin_attempt(now + 2).await,
ScannerPauseBacklogAttemptDecision::Tracked(_)
));
let decision = controller.begin_attempt(now + 2).await;
assert!(
matches!(decision, ScannerPauseBacklogAttemptDecision::Tracked(_)),
"expected tracked native expansion setup attempt, got {decision:?}; ledger={:?}; persistence_disabled={}; runtime_error={:?}",
controller.loaded.ledger,
controller.persistence_disabled,
runtime_error()
);
assert!(controller.loaded.ledger.has_unfinished_attempt());
}
let original = controller.loaded.ledger.clone();
@@ -1902,9 +1935,16 @@ mod tests {
old.pool_meta_write_status()
.await
.expect("healthy pool metadata keeps the background recovery loop read-only");
old.background_cancel_token().expect("old store shutdown token").cancel();
drop(old);
let expanded = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await;
shutdown_native_retirement_store(old).await;
let expanded = super::super::tests::setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root.path(),
false,
3,
2,
NATIVE_RETIREMENT_DRIVES_PER_SET,
false,
)
.await;
for pool in expanded.pools.iter().skip(old_pool_count) {
for set in &pool.disk_set {
let replica = read_scanner_pause_backlog_replica(Arc::clone(set)).await;
@@ -1999,12 +2039,16 @@ mod tests {
}
store.pool_meta_write_status().await.expect("healthy metadata before restart");
store
.background_cancel_token()
.expect("expanded store shutdown token")
.cancel();
drop(store);
let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await;
shutdown_native_retirement_store(store).await;
let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root.path(),
false,
3,
2,
NATIVE_RETIREMENT_DRIVES_PER_SET,
false,
)
.await;
let reloaded = load_scanner_pause_backlog(Arc::clone(&restarted))
.await
.expect("a new store must recover from disk without the failed controller");
@@ -2021,6 +2065,8 @@ mod tests {
assert_eq!(retry_ledger.current_attempt_serial, original.current_attempt_serial);
assert_eq!(retry_ledger.last_finished_attempt_serial, original.current_attempt_serial);
assert_eq!(retry_ledger.consecutive_failures, original.consecutive_failures + 1);
drop(retried);
shutdown_native_retirement_store(restarted).await;
}
async fn assert_current_native_writer_ledger(store: &Arc<ECStore>, expected: &ScannerPauseBacklogLedger) {
@@ -2124,6 +2170,8 @@ mod tests {
.await
.expect("retry must seed, commit and stabilize the stale member");
assert_current_native_writer_ledger(&store, &expected).await;
drop(target);
shutdown_native_retirement_store(store).await;
});
}
@@ -2175,6 +2223,8 @@ mod tests {
.await
.expect("a fresh caller may finish the seeded membership transition");
assert_current_native_writer_ledger(&store, &expected).await;
drop(pool_meta_lock);
shutdown_native_retirement_store(store).await;
});
}
@@ -2235,6 +2285,7 @@ mod tests {
.await
.expect("retry must seed the newly visible members before committing");
assert_current_native_writer_ledger(&store, &expected).await;
shutdown_native_retirement_store(store).await;
});
}
@@ -2292,6 +2343,7 @@ mod tests {
assert_current_native_ledger(&store, &original).await;
}
assert!(original.has_unfinished_attempt());
shutdown_native_retirement_store(store).await;
}
});
}
@@ -2337,14 +2389,22 @@ mod tests {
.pool_meta_write_status()
.await
.expect("healthy pool metadata keeps the background recovery loop read-only");
store.background_cancel_token().expect("old store shutdown token").cancel();
drop(store);
let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await;
shutdown_native_retirement_store(store).await;
let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root.path(),
false,
3,
2,
NATIVE_RETIREMENT_DRIVES_PER_SET,
false,
)
.await;
for set_index in 0..2 {
restarted.retire_scanner_pause_backlog_for_test(0, set_index).await.unwrap();
assert_native_source_missing(&restarted, set_index).await;
}
assert_current_native_ledger(&restarted, &original).await;
shutdown_native_retirement_store(restarted).await;
}
});
}
@@ -2384,9 +2444,16 @@ mod tests {
.pool_meta_write_status()
.await
.expect("healthy pool metadata keeps the background recovery loop read-only");
store.background_cancel_token().expect("old store shutdown token").cancel();
drop(store);
let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_sets(root.path(), false, 3, 2).await;
shutdown_native_retirement_store(store).await;
let restarted = super::super::tests::setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root.path(),
false,
3,
2,
NATIVE_RETIREMENT_DRIVES_PER_SET,
false,
)
.await;
for set_index in 0..2 {
restarted.retire_scanner_pause_backlog_for_test(0, set_index).await.unwrap();
assert_native_source_missing(&restarted, set_index).await;
@@ -2396,6 +2463,7 @@ mod tests {
assert_eq!(bootstrapped.ledger.generation, 1);
assert!(bootstrapped.ledger.last_updated_at_unix_secs < future_now);
assert_current_native_ledger(&restarted, &bootstrapped.ledger).await;
shutdown_native_retirement_store(restarted).await;
});
}
@@ -2441,6 +2509,8 @@ mod tests {
assert_current_native_ledger(&store, &original).await;
store.retire_scanner_pause_backlog_for_test(0, 1).await.unwrap();
assert_native_source_missing(&store, 1).await;
drop(pool_meta_lock);
shutdown_native_retirement_store(store).await;
});
}
@@ -2536,6 +2606,9 @@ mod tests {
"retirement never copied a stale source record"
);
}
drop(pool_meta_lock);
drop(barrier);
shutdown_native_retirement_store(store).await;
});
}
@@ -2669,11 +2742,12 @@ mod tests {
.retire_scanner_pause_backlog_for_test(0, 1)
.await
.expect("the remaining source set follows the same native proof");
let after = load_scanner_pause_backlog(store)
let after = load_scanner_pause_backlog(Arc::clone(&store))
.await
.expect("native restart selection after handoff");
assert_eq!(after.ledger, new_ledger);
assert!(after.durable && after.stable_matches_ledger);
shutdown_native_retirement_store(store).await;
});
}
@@ -2745,6 +2819,7 @@ mod tests {
.expect("retained target intent snapshot")
};
assert_eq!(after, before, "native cleanup never clears or estimates a target mutation intent");
shutdown_native_retirement_store(store).await;
});
}
+96 -5
View File
@@ -69,13 +69,42 @@ pub(super) async fn setup_scanner_cycle_store_at_path_with_sets(
seed_usage_baseline: bool,
pool_count: usize,
sets_per_pool: usize,
) -> Arc<ECStore> {
setup_scanner_cycle_store_at_path_with_layout(root, seed_usage_baseline, pool_count, sets_per_pool, 4).await
}
pub(super) async fn setup_scanner_cycle_store_at_path_with_layout(
root: &Path,
seed_usage_baseline: bool,
pool_count: usize,
sets_per_pool: usize,
drives_per_set: usize,
) -> Arc<ECStore> {
setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root,
seed_usage_baseline,
pool_count,
sets_per_pool,
drives_per_set,
true,
)
.await
}
pub(super) async fn setup_scanner_cycle_store_at_path_with_layout_and_disk_preinit(
root: &Path,
seed_usage_baseline: bool,
pool_count: usize,
sets_per_pool: usize,
drives_per_set: usize,
preinitialize_disks: bool,
) -> 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 set_index in 0..sets_per_pool {
for disk_index in 0..4 {
for disk_index in 0..drives_per_set {
let disk_path = if sets_per_pool == 1 {
root.join(format!("pool{pool_index}/disk{disk_index}"))
} else {
@@ -95,7 +124,7 @@ pub(super) async fn setup_scanner_cycle_store_at_path_with_sets(
pools.push(PoolEndpoints {
legacy: false,
set_count: sets_per_pool,
drives_per_set: 4,
drives_per_set,
endpoints: Endpoints::from(endpoints),
cmd_line: if pool_count == 1 && sets_per_pool == 1 {
"scanner-cycle-metrics".to_string()
@@ -108,9 +137,11 @@ pub(super) async fn setup_scanner_cycle_store_at_path_with_sets(
let endpoint_pools = EndpointServerPools::from(pools);
let instance_ctx = Arc::new(InstanceContext::new());
instance_ctx.set_endpoints(endpoint_pools.clone());
init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone())
.await
.expect("scanner cycle test disks should initialize");
if preinitialize_disks {
init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone())
.await
.expect("scanner cycle test disks should initialize");
}
let store = ECStore::new_with_instance_ctx(
"127.0.0.1:0".parse().expect("test address should parse"),
endpoint_pools,
@@ -7649,6 +7680,25 @@ async fn scanner_cycle_confirms_lost_remote_ack_from_activity_snapshot() {
"a new peer instance cannot confirm whether the old ACK reached durable dirty state"
);
let mut stale_activity = scanner_node_activity("epoch-a", 7, 3);
stale_activity.dirty_usage_generation = 4;
let stale_clean_activity = BTreeMap::from([("node-2".to_string(), stale_activity)]);
let stale_clean = remote_dirty_usage_acknowledgement_pending(
8,
1,
std::slice::from_ref(&acknowledgement),
std::future::ready(Err::<bool, _>(std::io::Error::other(
"response lost before newer generation was observed",
))),
|| async { Ok(stale_clean_activity) },
)
.await;
assert_eq!(
scanner_cycle_outcome_with_pending_maintenance(ScannerCycleOutcome::Completed, stale_clean),
ScannerCycleOutcome::CompletedWithPendingMaintenance,
"a clean peer snapshot from before the acknowledged generation cannot prove the ACK reached durable dirty state"
);
let mut written_activity = scanner_node_activity("epoch-a", 7, 3);
written_activity.dirty_usage_generation = 6;
written_activity.dirty_usage_pending = true;
@@ -7722,6 +7772,47 @@ async fn scanner_cycle_confirms_lost_scoped_ack_only_after_same_instance_clean_a
"a restarted peer cannot prove the scoped ACK reached the old scanner instance"
);
let mut stale_activity = scanner_node_activity("epoch-a", 7, 3);
stale_activity.dirty_usage_generation = 4;
let stale_clean_activity = BTreeMap::from([("node-2".to_string(), stale_activity)]);
let stale_clean = remote_dirty_usage_acknowledgement_pending(
8,
1,
std::slice::from_ref(&acknowledgement),
std::future::ready(Err::<bool, _>(std::io::Error::other(
"scoped ACK transport failed before the requested generation was observed",
))),
|| async { Ok(stale_clean_activity) },
)
.await;
assert_eq!(
scanner_cycle_outcome_with_pending_maintenance(ScannerCycleOutcome::Completed, stale_clean),
ScannerCycleOutcome::CompletedWithPendingMaintenance,
"a clean peer snapshot from before the scoped ACK generation cannot prove the ACK reached durable dirty state"
);
let empty_scoped_ack = ScannerDirtyUsageAcknowledgement {
host: "node-2".to_string(),
instance_id: "epoch-a".to_string(),
kind: ScannerDirtyUsageAcknowledgementKind::Scoped {
owner_id: Uuid::from_u128(0x11111111111111111111111111111111).to_string(),
entries: Vec::new(),
},
};
let empty_scoped_clean = remote_dirty_usage_acknowledgement_pending(
8,
1,
&[empty_scoped_ack],
std::future::ready(Err::<bool, _>(std::io::Error::other("empty scoped ACK failed before peer delivery"))),
|| async { Ok(BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))])) },
)
.await;
assert_eq!(
scanner_cycle_outcome_with_pending_maintenance(ScannerCycleOutcome::Completed, empty_scoped_clean),
ScannerCycleOutcome::CompletedWithPendingMaintenance,
"an empty scoped ACK has no durable generation to reconcile after response loss"
);
let mut written_activity = scanner_node_activity("epoch-a", 7, 3);
written_activity.dirty_usage_generation = 6;
written_activity.dirty_usage_pending = true;
@@ -432,7 +432,12 @@ async fn scoped_ack_publication_stale_baseline_cannot_prove_a_replaced_root() {
#[tokio::test]
#[serial]
async fn scoped_ack_publication_rejects_builder_mutation_after_real_root_publish() {
for mutation in ["remote_ack_target", "publication_epoch", "remote_lease_targets"] {
for mutation in [
"remote_ack_target",
"remote_scoped_ack_target",
"publication_epoch",
"remote_lease_targets",
] {
let (_directory, store) = candidate_store().await;
let (scan, candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let baseline = read_data_usage_persist_baseline(store.clone())
@@ -465,6 +470,18 @@ async fn scoped_ack_publication_rejects_builder_mutation_after_real_root_publish
instance_id: crate::scanner_activity_epoch().to_string(),
kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Generation(changed_generation),
}]),
"remote_scoped_ack_target" => scan.with_remote_dirty_usage_acknowledgements(vec![ScannerDirtyUsageAcknowledgement {
host: "proof-peer:9000".to_string(),
instance_id: crate::scanner_activity_epoch().to_string(),
kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped {
owner_id: crate::scanner_activity_epoch().to_string(),
entries: vec![crate::storage_api::EcstoreScannerScopedDirtyUsageAckEntry {
bucket: PROOF_BUCKET.to_string(),
bucket_incarnation: uuid::Uuid::from_u128(0x11111111111111111111111111111111),
generation: changed_generation,
}],
},
}]),
"publication_epoch" => scan.with_publication_epoch(Some(changed_epoch)),
"remote_lease_targets" => scan.with_remote_publication_lease_targets(vec![(
"proof-peer:9000".to_string(),
+62 -18
View File
@@ -306,7 +306,7 @@ fn scanner_scoped_dirty_usage_ack_exceeds_cost_threshold(
fn resolve_remote_dirty_usage_scope(
requested_scope: ScannerBucketScanScope,
mut dirty_buckets: HashSet<String>,
dirty_usage_snapshot: &DirtyUsageSnapshot,
remote_dirty_usage: VerifiedRemoteDirtyUsage,
all_buckets: &[BucketInfo],
baseline_proof: ScannerCacheBaselineProof<'_>,
@@ -319,12 +319,21 @@ fn resolve_remote_dirty_usage_scope(
let peer_count = remote_dirty_usage.peer_count;
let dirty_peer_count = remote_dirty_usage.dirty_peer_count;
let remote_dirty_buckets = remote_dirty_usage.dirty_buckets.clone();
let mut dirty_buckets = dirty_usage_snapshot.buckets.keys().cloned().collect::<HashSet<_>>();
dirty_buckets.extend(remote_dirty_usage.dirty_buckets);
// Peer snapshots contribute bucket names only; the local prefix scopes
// would narrow a bucket a peer dirtied elsewhere, so the merged scope
// stays at bucket granularity (same rule as the local fallthrough).
let scope =
scoped_scan_scope_from_dirty_buckets(requested_scope, dirty_buckets, None, true, false, all_buckets, baseline_proof);
let scope = scoped_scan_scope_from_dirty_buckets(
requested_scope.clone(),
dirty_buckets.clone(),
None,
true,
false,
all_buckets,
baseline_proof,
);
if scope.is_default() {
return default_result(scope);
}
@@ -358,19 +367,45 @@ fn resolve_remote_dirty_usage_scope(
}
let has_scoped_acknowledgements = !scoped_acknowledgements.is_empty();
let distributed_segment_invalidation_evidence =
(dirty_peer_count > 0 && has_scoped_acknowledgements).then_some(DistributedSegmentInvalidationEvidence {
invalidation_domain: crate::segment_invalidation::SegmentInvalidationDomain::DistributedEc,
distributed_ec_invalidation: true,
peer_count,
dirty_peer_count,
same_window_remote_proof: true,
all_peers_bound_to_generation_window: true,
});
let segment_reuse_activation_preflight = scanner_segment_reuse_activation_preflight_for_baseline_with_evidence(
dirty_usage_snapshot,
true,
baseline_proof,
distributed_segment_invalidation_evidence,
);
let scope = if segment_reuse_activation_preflight.scanner_segment_reuse_activated {
let local_only_scopes = dirty_usage_snapshot
.scopes
.iter()
.filter(|(bucket, _)| !remote_dirty_buckets.contains(bucket.as_str()))
.map(|(bucket, scope)| (bucket.clone(), scope.clone()))
.collect::<DirtyUsageBucketScopes>();
scoped_scan_scope_from_dirty_buckets(
requested_scope,
dirty_buckets,
(!local_only_scopes.is_empty()).then_some(&local_only_scopes),
true,
true,
all_buckets,
baseline_proof,
)
} else {
scope
};
ScannerBucketScopeResolutionResult {
scope,
remote_dirty_usage_acknowledgements: scoped_acknowledgements,
distributed_segment_invalidation_evidence: (dirty_peer_count > 0 && has_scoped_acknowledgements).then_some(
DistributedSegmentInvalidationEvidence {
invalidation_domain: crate::segment_invalidation::SegmentInvalidationDomain::DistributedEc,
distributed_ec_invalidation: true,
peer_count,
dirty_peer_count,
same_window_remote_proof: true,
all_peers_bound_to_generation_window: true,
},
),
distributed_segment_invalidation_evidence,
}
}
@@ -596,7 +631,7 @@ fn scanner_segment_reuse_activation_preflight_for_cycle(
production_activation: true,
producer_identity_coverage_complete: dirty_usage_producer_evidence.producer_identity_coverage_complete,
durable_producer_identity: dirty_usage_producer_evidence.durable_producer_identity,
durable_dirty_producer_journal: dirty_usage_producer_evidence.durable_producer_identity,
durable_dirty_producer_journal: dirty_usage_producer_evidence.durable_dirty_producer_journal,
restart_gap_absent: dirty_usage_producer_evidence.restart_gap_absent,
generation_window_bound: dirty_usage_snapshot.covers_all_pending
&& dirty_usage_snapshot.generation != 0
@@ -616,6 +651,15 @@ fn scanner_segment_reuse_activation_preflight_for_baseline(
dirty_usage_snapshot: &DirtyUsageSnapshot,
distributed: bool,
baseline_proof: ScannerCacheBaselineProof<'_>,
) -> ScannerSegmentReuseActivationPreflight {
scanner_segment_reuse_activation_preflight_for_baseline_with_evidence(dirty_usage_snapshot, distributed, baseline_proof, None)
}
fn scanner_segment_reuse_activation_preflight_for_baseline_with_evidence(
dirty_usage_snapshot: &DirtyUsageSnapshot,
distributed: bool,
baseline_proof: ScannerCacheBaselineProof<'_>,
distributed_segment_invalidation_evidence: Option<DistributedSegmentInvalidationEvidence>,
) -> ScannerSegmentReuseActivationPreflight {
let (dirty_usage_producer_evidence, cold_zero_walk_oracle) =
scanner_segment_reuse_baseline_producer_evidence(dirty_usage_snapshot, baseline_proof);
@@ -623,7 +667,7 @@ fn scanner_segment_reuse_activation_preflight_for_baseline(
dirty_usage_snapshot,
dirty_usage_producer_evidence,
distributed,
None,
distributed_segment_invalidation_evidence,
cold_zero_walk_oracle,
)
}
@@ -1568,14 +1612,14 @@ pub(crate) use cache::{
current_cache_root_or_prepare_with_generation,
};
pub use dirty_usage::{
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageClearObserver, ScannerDirtyUsageSnapshot,
ScannerDirtyUsageState, ScannerDurableDirtyUsageReplayEntry, ScannerDurableDirtyUsageReplayError,
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageClearObserver, ScannerDirtyUsageMutationObserver,
ScannerDirtyUsageSnapshot, ScannerDirtyUsageState, ScannerDurableDirtyUsageReplayEntry, ScannerDurableDirtyUsageReplayError,
ScannerDurableDirtyUsageReplayRecord, ScannerDurableDirtyUsageReplayScope, acknowledge_dirty_usage_generation,
acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, encode_durable_dirty_usage_producer_replay_record,
record_dirty_usage_bucket, record_dirty_usage_bucket_from_producer, record_dirty_usage_bucket_from_producers,
record_dirty_usage_object, record_dirty_usage_object_from_producer, record_scanner_maintenance_change,
replay_durable_dirty_usage_producer_record, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state,
scanner_maintenance_generation, set_scanner_dirty_usage_clear_observer,
scanner_maintenance_generation, set_scanner_dirty_usage_clear_observer, set_scanner_dirty_usage_mutation_observer,
};
#[cfg(test)]
pub(crate) use dirty_usage::{clear_dirty_usage_buckets_for_tests, dirty_usage_buckets_for_tests};
+93 -5
View File
@@ -30,6 +30,8 @@ pub(super) static DIRTY_USAGE_PRODUCER_IDENTITIES: LazyLock<StdMutex<DirtyUsageP
LazyLock::new(|| StdMutex::new(BTreeMap::new()));
static DIRTY_USAGE_CLEAR_OBSERVER: LazyLock<StdRwLock<Option<ScannerDirtyUsageClearObserver>>> =
LazyLock::new(|| StdRwLock::new(None));
static DIRTY_USAGE_MUTATION_OBSERVER: LazyLock<StdRwLock<Option<ScannerDirtyUsageMutationObserver>>> =
LazyLock::new(|| StdRwLock::new(None));
pub(super) static DIRTY_USAGE_PRODUCER_COVERAGE: AtomicU64 = AtomicU64::new(0);
pub(super) static DIRTY_USAGE_BUCKET_NOTIFY: LazyLock<Notify> = LazyLock::new(Notify::new);
pub(super) static SCANNER_ACTIVITY_EPOCH: LazyLock<String> = LazyLock::new(|| format!("{:032x}", rand::random::<u128>()));
@@ -53,6 +55,8 @@ pub struct ScannerDirtyUsageBucket {
}
pub type ScannerDirtyUsageClearObserver = Arc<dyn Fn(Vec<ScannerDirtyUsageBucket>) + Send + Sync + 'static>;
pub type ScannerDirtyUsageMutationObserver =
Arc<dyn Fn(&str, &str, crate::segment_invalidation::SegmentInvalidationProducerIdentity) + Send + Sync + 'static>;
/// A non-durable optimization hint for a dirty bucket.
///
@@ -83,6 +87,7 @@ pub(super) type DirtyUsageProducerIdentities = BTreeMap<String, DirtyUsageProduc
pub(super) struct DirtyUsageProducerEvidence {
pub(super) producer_identity_coverage_complete: bool,
pub(super) durable_producer_identity: bool,
pub(super) durable_dirty_producer_journal: bool,
pub(super) restart_gap_absent: bool,
pub(super) generation_window_bound: bool,
pub(super) generation_start: u64,
@@ -112,6 +117,15 @@ pub fn set_scanner_dirty_usage_clear_observer(
std::mem::replace(&mut *slot, observer)
}
pub fn set_scanner_dirty_usage_mutation_observer(
observer: Option<ScannerDirtyUsageMutationObserver>,
) -> Option<ScannerDirtyUsageMutationObserver> {
let mut slot = DIRTY_USAGE_MUTATION_OBSERVER
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner());
std::mem::replace(&mut *slot, observer)
}
fn notify_dirty_usage_clear(cleared: Vec<ScannerDirtyUsageBucket>) {
if cleared.is_empty() {
return;
@@ -125,6 +139,30 @@ fn notify_dirty_usage_clear(cleared: Vec<ScannerDirtyUsageBucket>) {
}
}
fn notify_dirty_usage_mutation(
bucket: &str,
object: &str,
producers: &[crate::segment_invalidation::SegmentInvalidationProducerIdentity],
) {
if bucket.is_empty() {
return;
}
let observer = DIRTY_USAGE_MUTATION_OBSERVER
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clone();
let Some(observer) = observer else {
return;
};
if producers.is_empty() {
observer(bucket, object, crate::segment_invalidation::SegmentInvalidationProducerIdentity::Unknown);
} else {
for producer in producers {
observer(bucket, object, *producer);
}
}
}
/// A point-in-time view of the local dirty bucket generations.
///
/// `complete == false` is an all-or-nothing overflow signal: `buckets` is
@@ -566,6 +604,7 @@ mod scoped_dirty_usage_tests {
let evidence = dirty_usage_producer_evidence(&snapshot);
assert!(evidence.producer_identity_coverage_complete);
assert!(evidence.durable_producer_identity);
assert!(evidence.durable_dirty_producer_journal);
assert!(evidence.restart_gap_absent);
assert_eq!(evidence.generation_start, 7);
assert_eq!(evidence.generation_end, 7);
@@ -583,10 +622,45 @@ mod scoped_dirty_usage_tests {
);
let changed_evidence = dirty_usage_producer_evidence(&changed);
assert!(!changed_evidence.durable_producer_identity);
assert!(!changed_evidence.durable_dirty_producer_journal);
assert!(!changed_evidence.restart_gap_absent);
clear_dirty_usage_buckets_for_tests();
}
#[test]
#[serial]
fn dirty_usage_mutation_observer_tracks_typed_and_conservative_events() {
use std::sync::{Arc, Mutex};
clear_dirty_usage_buckets_for_tests();
let observed = Arc::new(Mutex::new(Vec::new()));
let observed_clone = observed.clone();
let previous = set_scanner_dirty_usage_mutation_observer(Some(Arc::new(move |bucket, object, producer| {
observed_clone.lock().expect("observer lock should not be poisoned").push((
bucket.to_string(),
object.to_string(),
producer,
));
})));
record_dirty_usage_object_from_producer("photos", "2026/image.jpg", SegmentInvalidationProducerIdentity::PutObject);
record_dirty_usage_bucket("archive");
set_scanner_dirty_usage_mutation_observer(previous);
assert_eq!(
*observed.lock().expect("observer lock should not be poisoned"),
vec![
(
"photos".to_string(),
"2026/image.jpg".to_string(),
SegmentInvalidationProducerIdentity::PutObject,
),
("archive".to_string(), String::new(), SegmentInvalidationProducerIdentity::Unknown,),
]
);
clear_dirty_usage_buckets_for_tests();
}
#[test]
#[serial]
fn durable_dirty_usage_replay_rejects_invalid_records_without_partial_state() {
@@ -745,6 +819,7 @@ fn record_dirty_usage_bucket_inner<I>(bucket: &str, producers: I)
where
I: IntoIterator<Item = crate::segment_invalidation::SegmentInvalidationProducerIdentity>,
{
let producers = producers.into_iter().collect::<Vec<_>>();
let pending_buckets = {
let mut dirty_buckets = dirty_usage_buckets();
let mut dirty_scopes = dirty_usage_bucket_scopes();
@@ -752,7 +827,12 @@ where
let generation = advance_generation(&DIRTY_USAGE_BUCKET_GENERATION);
dirty_buckets.insert(bucket.to_string(), generation);
dirty_scopes.insert(bucket.to_string(), DirtyUsageBucketScope::WholeBucket);
record_segment_invalidation_producer_identities_for_generation(&mut producer_identities, bucket, generation, producers);
record_segment_invalidation_producer_identities_for_generation(
&mut producer_identities,
bucket,
generation,
producers.iter().copied(),
);
dirty_buckets.len()
};
global_metrics().record_scanner_dirty_usage_pending(usize_to_u64_saturated(pending_buckets));
@@ -760,6 +840,7 @@ where
// admin/console consumers never ride the full TTL after a change
// (rustfs/backlog#1872).
crate::prefix_usage::invalidate_prefix_usage_cache(bucket);
notify_dirty_usage_mutation(bucket, "", &producers);
DIRTY_USAGE_BUCKET_NOTIFY.notify_one();
}
@@ -817,11 +898,17 @@ where
if overflowed {
*scope = DirtyUsageBucketScope::WholeBucket;
}
record_segment_invalidation_producer_identities_for_generation(&mut producer_identities, bucket, generation, producers);
record_segment_invalidation_producer_identities_for_generation(
&mut producer_identities,
bucket,
generation,
producers.iter().copied(),
);
dirty_buckets.len()
};
global_metrics().record_scanner_dirty_usage_pending(usize_to_u64_saturated(pending_buckets));
crate::prefix_usage::invalidate_prefix_usage_cache(bucket);
notify_dirty_usage_mutation(bucket, object, &producers);
DIRTY_USAGE_BUCKET_NOTIFY.notify_one();
}
@@ -1193,7 +1280,7 @@ pub(super) fn dirty_usage_producer_evidence(snapshot: &DirtyUsageSnapshot) -> Di
&& DIRTY_USAGE_PRODUCER_COVERAGE.load(Ordering::Acquire)
& crate::segment_invalidation::SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION_COVERAGE_MASK
== crate::segment_invalidation::SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION_COVERAGE_MASK;
let durable_producer_identity = producer_identity_coverage_complete
let durable_dirty_producer_journal = producer_identity_coverage_complete
&& snapshot
.buckets
.keys()
@@ -1205,8 +1292,9 @@ pub(super) fn dirty_usage_producer_evidence(snapshot: &DirtyUsageSnapshot) -> Di
DirtyUsageProducerEvidence {
producer_identity_coverage_complete,
durable_producer_identity,
restart_gap_absent: durable_producer_identity,
durable_producer_identity: durable_dirty_producer_journal,
durable_dirty_producer_journal,
restart_gap_absent: durable_dirty_producer_journal,
generation_window_bound,
generation_start: if generation_window_bound { generation_start } else { 0 },
generation_end: if generation_window_bound { generation_end } else { 0 },
+1 -1
View File
@@ -190,7 +190,7 @@ where
};
let remote_resolution = resolve_remote_dirty_usage_scope(
resolution.requested_scope,
dirty_buckets,
resolution.dirty_usage_snapshot,
remote_dirty_usage,
resolution.all_buckets,
resolution.baseline_proof,
+151 -5
View File
@@ -37,7 +37,7 @@ use rustfs_concurrency::{
};
use rustfs_filemeta::FileInfo;
use serial_test::serial;
use std::collections::BTreeMap;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use temp_env::with_var;
use time::OffsetDateTime;
@@ -285,6 +285,7 @@ fn scanner_durable_segment_invalidation_evidence_requires_matching_complete_set_
assert!(durable_evidence.producer_identity_coverage_complete);
assert!(durable_evidence.durable_producer_identity);
assert!(!durable_evidence.durable_dirty_producer_journal);
assert!(durable_evidence.restart_gap_absent);
let mut stale_epoch = results.clone();
@@ -297,12 +298,14 @@ fn scanner_durable_segment_invalidation_evidence_requires_matching_complete_set_
let stale_evidence = scanner_durable_segment_invalidation_evidence(&dirty_usage_snapshot, &stale_epoch, &expected_sources);
assert!(stale_evidence.producer_identity_coverage_complete);
assert!(!stale_evidence.durable_producer_identity);
assert!(!stale_evidence.durable_dirty_producer_journal);
assert!(!stale_evidence.restart_gap_absent);
record_dirty_usage_bucket("videos");
let changed_evidence = scanner_durable_segment_invalidation_evidence(&dirty_usage_snapshot, &results, &expected_sources);
assert!(!changed_evidence.producer_identity_coverage_complete);
assert!(!changed_evidence.durable_producer_identity);
assert!(!changed_evidence.durable_dirty_producer_journal);
assert!(!changed_evidence.restart_gap_absent);
clear_dirty_usage_buckets_for_tests();
}
@@ -316,6 +319,19 @@ fn scanner_segment_reuse_activation_replays_cold_durable_baseline() {
for producer in SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION {
record_dirty_usage_object_from_producer("photos", "2026/object", producer);
}
let replay_generation = dirty_usage_generation();
replay_durable_dirty_usage_producer_record(
&encode_durable_dirty_usage_producer_replay_record(vec![ScannerDurableDirtyUsageReplayEntry {
bucket: "photos".to_string(),
generation: replay_generation,
scope: ScannerDurableDirtyUsageReplayScope::TopLevelEntries {
entries: BTreeSet::from(["2026".to_string()]),
},
producers: SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION.into_iter().collect(),
}])
.expect("durable producer replay should encode"),
)
.expect("durable producer replay should restore restart authority");
let dirty_usage_snapshot =
snapshot_dirty_usage_buckets(&[bucket_info("photos"), bucket_info("archive")], dirty_usage_generation());
let mut segment_proof = dirty_usage_producer_evidence(&dirty_usage_snapshot)
@@ -413,8 +429,11 @@ fn scanner_segment_reuse_activation_replays_cold_durable_baseline() {
scan_plan_digest,
},
);
assert!(preflight.scanner_segment_reuse_activated);
assert_eq!(preflight.fail_closed_blockers().collect::<Vec<_>>(), Vec::<&str>::new());
assert!(!preflight.scanner_segment_reuse_activated);
assert_eq!(
preflight.fail_closed_blockers().collect::<Vec<_>>(),
vec!["missing_durable_journal_replay"]
);
record_dirty_usage_bucket("photos");
let unidentified_snapshot =
@@ -473,6 +492,7 @@ fn complete_process_local_producer_evidence() -> DirtyUsageProducerEvidence {
DirtyUsageProducerEvidence {
producer_identity_coverage_complete: true,
durable_producer_identity: false,
durable_dirty_producer_journal: false,
restart_gap_absent: false,
generation_window_bound: true,
generation_start: 7,
@@ -1468,6 +1488,7 @@ fn dirty_usage_producer_evidence_tracks_process_local_coverage_without_durable_r
assert!(evidence.generation_window_bound);
assert!(evidence.producer_identity_coverage_complete);
assert!(!evidence.durable_producer_identity);
assert!(!evidence.durable_dirty_producer_journal);
assert!(!evidence.restart_gap_absent);
assert_eq!(evidence.generation_start, snapshot.buckets["photos"]);
assert_eq!(evidence.generation_end, snapshot.buckets["photos"]);
@@ -1858,6 +1879,44 @@ fn complete_usage_baseline(
bytes::Bytes::from(serde_json::to_vec(&baseline).expect("test baseline should encode"))
}
fn complete_segment_reuse_baseline(
source: DataUsageCacheSource,
scan_plan_digest: DataUsageScanPlanDigest,
scanner_cycle: u64,
scanner_epoch: u64,
evidence: DirtyUsageProducerEvidence,
buckets: &[&str],
) -> bytes::Bytes {
let mut proof = evidence
.segment_invalidation_proof()
.expect("durable producer evidence should produce segment proof");
proof.cold_zero_walk_oracle = true;
let baseline = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
scanner_cycle: Some(scanner_cycle),
scanner_epoch: Some(scanner_epoch),
buckets_count: u64::try_from(buckets.len()).expect("test bucket count should fit"),
buckets_usage: buckets
.iter()
.map(|bucket| ((*bucket).to_string(), Default::default()))
.collect(),
usage_snapshot_complete: true,
usage_snapshot_converged: Some(true),
usage_snapshot_set_states: vec![DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).expect("test pool index should fit"),
set_index: u64::try_from(source.set_index).expect("test set index should fit"),
scanner_cycle: Some(scanner_cycle),
scanner_epoch: Some(scanner_epoch),
scan_plan_digest: Some(scan_plan_digest.0),
complete: true,
tombstone: false,
segment_invalidation_proof: Some(proof),
}],
..Default::default()
};
bytes::Bytes::from(serde_json::to_vec(&baseline).expect("test baseline should encode"))
}
#[test]
fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance() {
let source = DataUsageCacheSource::new(1, 2);
@@ -2183,6 +2242,12 @@ fn remote_dirty_usage_invalidates_local_prefix_hints_until_distributed_proof_exi
"photos".to_string(),
DirtyUsageBucketScope::TopLevelEntries(HashSet::from(["2026".to_string()])),
)]);
let dirty_usage_snapshot = DirtyUsageSnapshot {
buckets: Arc::new(HashMap::from([("photos".to_string(), 7)])),
scopes: Arc::new(dirty_scopes.clone()),
generation: 7,
covers_all_pending: true,
};
let locally_scoped = scoped_scan_scope_from_dirty_buckets(
ScannerBucketScanScope::default(),
HashSet::from(["photos".to_string()]),
@@ -2206,7 +2271,7 @@ fn remote_dirty_usage_invalidates_local_prefix_hints_until_distributed_proof_exi
let distributed = resolve_remote_dirty_usage_scope(
ScannerBucketScanScope::default(),
HashSet::from(["photos".to_string()]),
&dirty_usage_snapshot,
remote_dirty_usage,
&[bucket_info("photos")],
ScannerCacheBaselineProof {
@@ -2246,6 +2311,82 @@ fn remote_dirty_usage_invalidates_local_prefix_hints_until_distributed_proof_exi
assert!(evidence.all_peers_bound_to_generation_window);
}
#[test]
#[serial]
fn distributed_segment_reuse_activation_keeps_remote_dirty_buckets_at_bucket_scope() {
clear_dirty_usage_buckets_for_tests();
let source = DataUsageCacheSource::new(1, 2);
let expected_sources = HashSet::from([source]);
let scan_plan_digest = DataUsageScanPlanDigest([9; 32]);
let entries = BTreeSet::from(["2026".to_string()]);
let bytes = encode_durable_dirty_usage_producer_replay_record(vec![ScannerDurableDirtyUsageReplayEntry {
bucket: "photos".to_string(),
generation: 7,
scope: ScannerDurableDirtyUsageReplayScope::TopLevelEntries { entries },
producers: crate::segment_invalidation::SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION
.iter()
.copied()
.collect(),
}])
.expect("durable dirty usage replay record should encode");
replay_durable_dirty_usage_producer_record(&bytes).expect("durable dirty usage replay should restore producer proof");
let dirty_usage_snapshot =
snapshot_dirty_usage_buckets(&[bucket_info("photos"), bucket_info("archive")], dirty_usage_generation());
let evidence = dirty_usage_producer_evidence(&dirty_usage_snapshot);
let baseline = complete_segment_reuse_baseline(source, scan_plan_digest, 7, 11, evidence, &["photos", "archive"]);
let expected_peers = HashMap::from([(
"node-a:9000".to_string(),
ScannerPeerDirtyUsageExpectation {
instance_id: "instance-a".to_string(),
generation: 3,
pending: true,
},
)]);
let remote_dirty_usage = verified_remote_dirty_usage(
&expected_peers,
vec![(
"node-a:9000".to_string(),
peer_dirty_usage_snapshot("instance-a", 3, true, &[("archive", 3)]),
)],
)
.expect("fixture remote dirty usage should verify at bucket granularity");
let result = resolve_remote_dirty_usage_scope(
ScannerBucketScanScope::default(),
&dirty_usage_snapshot,
remote_dirty_usage,
&[bucket_info("photos"), bucket_info("archive")],
ScannerCacheBaselineProof {
authoritative_data: Some(&baseline),
observed_candidate_data: None,
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
scan_plan_digest,
},
);
assert_eq!(
result
.scope
.selected_buckets
.as_deref()
.expect("distributed reuse still selects both dirty buckets"),
&HashSet::from(["photos".to_string(), "archive".to_string()])
);
assert!(
result.scope.prefix_scope_for("photos").is_some(),
"durable local producer proof may activate local segment reuse after distributed ACK capability evidence"
);
assert!(
result.scope.prefix_scope_for("archive").is_none(),
"remote dirty usage has only bucket-granularity evidence and must not be narrowed to a local prefix"
);
assert_eq!(result.remote_dirty_usage_acknowledgements.len(), 1);
assert!(result.distributed_segment_invalidation_evidence.is_some());
clear_dirty_usage_buckets_for_tests();
}
fn peer_dirty_usage_snapshot(
instance_id: &str,
generation: u64,
@@ -2419,7 +2560,12 @@ fn remote_dirty_usage_scope_resolution_falls_back_when_ack_batch_exceeds_thresho
let result = resolve_remote_dirty_usage_scope(
ScannerBucketScanScope::default(),
HashSet::new(),
&DirtyUsageSnapshot {
buckets: Arc::new(HashMap::new()),
scopes: Arc::new(HashMap::new()),
generation: 7,
covers_all_pending: true,
},
remote_dirty_usage,
&all_buckets,
ScannerCacheBaselineProof {
@@ -226,6 +226,34 @@ fn record_segment_dirty_usage(bucket: &str) {
}
}
fn replay_segment_dirty_usage(bucket: &str) {
replay_dirty_usage(
bucket,
ScannerDurableDirtyUsageReplayScope::TopLevelEntries {
entries: BTreeSet::from(["hot-segment".to_string()]),
},
);
}
fn replay_whole_bucket_dirty_usage(bucket: &str) {
replay_dirty_usage(bucket, ScannerDurableDirtyUsageReplayScope::WholeBucket);
}
fn replay_dirty_usage(bucket: &str, scope: ScannerDurableDirtyUsageReplayScope) {
replay_durable_dirty_usage_producer_record(
&encode_durable_dirty_usage_producer_replay_record(vec![ScannerDurableDirtyUsageReplayEntry {
bucket: bucket.to_string(),
generation: dirty_usage_generation(),
scope,
producers: crate::segment_invalidation::SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION
.into_iter()
.collect(),
}])
.expect("durable segment replay should encode"),
)
.expect("durable segment replay should restore producer authority");
}
// The scoped fallback fixture keeps two EC pools and several scan futures live
// at once. Run the async cases on a dedicated stack so Linux libtest defaults
// exercise the assertions instead of aborting before the oracle finishes.
@@ -267,6 +295,7 @@ async fn scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks_
create_bucket(&store, &hot).await;
create_bucket(&store, &cold).await;
record_segment_dirty_usage(&hot);
replay_segment_dirty_usage(&hot);
let baseline = run_entry(&store, 1, None, true, false, false).await;
persist_baseline(&store, &baseline).await;
@@ -283,6 +312,7 @@ async fn scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks_
"hot-segment/object",
crate::segment_invalidation::SegmentInvalidationProducerIdentity::PutObject,
);
replay_segment_dirty_usage(&hot);
let usage = run_entry(&store, 3, Some(&hot), true, true, true).await;
assert_eq!(usage.buckets_usage[&hot].objects_count, 2);
assert_eq!(usage.buckets_usage[&cold].objects_count, 1);
@@ -298,6 +328,7 @@ async fn scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks_
crate::segment_invalidation::SegmentInvalidationProducerIdentity::PutObject,
);
}
replay_whole_bucket_dirty_usage(&hot);
let usage = run_entry(&store, 4, Some(&hot), true, true, false).await;
assert_eq!(usage.objects_total_count, 3);
+7 -5
View File
@@ -416,12 +416,14 @@ sample, or save-frequency cost counters. Missing, synthetic, stale, tampered,
undersized, cross-run, or topology-mismatched evidence returns a compact blocked
or invalid JSON result and a nonzero exit.
The status-and-outcome descriptor producer consumes three measured raw JSON
artifacts for G05, G06, and R-D. Those inputs must all carry schema 1, measured
evidence, matching `source_revision`, a shared `run_id`, a shared
The status-and-outcome raw collector normalizes live observations into the three
measured raw JSON artifacts for G05, G06, and R-D. The descriptor producer then
consumes those artifacts. The inputs must all carry schema 1, measured evidence,
matching `source_revision`, a shared `run_id`, a shared
`measurement_window_id`, matching `started_at`/`finished_at` timestamps, and
non-empty command provenance. The producer rejects command-line run/window/time
overrides that would relabel raw artifacts from another status-and-outcome run.
non-empty command provenance. The collector and producer reject synthetic input,
missing required cases, and command-line run/window/time overrides that would
relabel raw artifacts from another status-and-outcome run.
The scheduler-pressure lane must also carry the numbers needed to close W09,
W10, and W11: bounded deferred item/byte/age limits, zero duplicate tasks,
+5 -2
View File
@@ -5939,8 +5939,11 @@ mod tests {
})
.await
.expect_err("mismatched kms context should fail");
assert_eq!(err.code, S3ErrorCode::InternalError);
assert_eq!(err.message, ApiError::error_code_to_message(&S3ErrorCode::InternalError));
assert_eq!(err.code, S3ErrorCode::InvalidRequest);
assert_eq!(
err.message,
"Encryption context mismatch: Context mismatch for key 'tenant': expected 'alpha', got 'beta'"
);
assert_eq!(super::kms_data_plane_error_class(&err), "context_mismatch");
manager.stop().await.expect("kms service should stop cleanly");
+4 -2
View File
@@ -955,7 +955,10 @@ pub(crate) async fn get_local_server_property() -> rustfs_madmin::ServerProperti
pub(crate) async fn init_background_replication(store: Arc<ECStore>) {
let durable_dirty_usage_journal = super::scanner_dirty_journal::start_durable_dirty_usage_journal(store.clone()).await;
let journal_writer = durable_dirty_usage_journal.clone();
let mutation_journal = durable_dirty_usage_journal.clone();
rustfs_scanner::set_scanner_dirty_usage_mutation_observer(Some(Arc::new(move |bucket, object, producer| {
mutation_journal.record_committed_mutation(bucket, object, producer);
})));
ecstore_bucket::replication::set_scanner_dirty_usage_mutation_observer(Some(Arc::new(move |bucket, object, source| {
let producer = match source {
ecstore_bucket::replication::ScannerDirtyUsageMutationSource::Replication => {
@@ -966,7 +969,6 @@ pub(crate) async fn init_background_replication(store: Arc<ECStore>) {
}
};
rustfs_scanner::record_dirty_usage_object_from_producer(bucket, object, producer);
journal_writer.record_committed_mutation(bucket, object, producer);
})));
rustfs_scanner::set_scanner_dirty_usage_clear_observer(Some(Arc::new(move |cleared| {
durable_dirty_usage_journal.clear_confirmed_buckets(cleared);
+1
View File
@@ -65,6 +65,7 @@ their issue closes.
| `run_scanner_heal_legacy_rollback_evidence.py` | dev-tool | Assembles measured Scanner/Heal R-L legacy source-conflict, migration-gap, and source-retirement release descriptors | `.config/scanner-heal-required-tests.json`; `test_scanner_heal_legacy_rollback_evidence.sh` |
| `run_scanner_heal_mrf_evidence.py` | dev-tool | Assembles measured Scanner/Heal G07/G08/P4 MRF release descriptors from W13 raw artifacts | `.config/scanner-heal-required-tests.json`; `test_scanner_heal_w13_mrf_evidence.sh` |
| `run_scanner_heal_scheduler_pressure_evidence.py` | dev-tool | Assembles measured Scanner/Heal G10/P1/P3 scheduler-pressure release descriptors from a completed measured ABBA run, recovery-window proof, and profile artifacts | `docs/operations/scanner-benchmark-runbook.md`; `test_scanner_heal_scheduler_pressure_evidence.sh` |
| `run_scanner_heal_status_outcome_probe.py` | dev-tool | Normalizes live Scanner/Heal status/outcome observations into the measured G05/G06/R-D raw artifacts consumed by the descriptor producer | `.config/scanner-heal-required-tests.json`; `test_scanner_heal_status_outcome_evidence.sh` |
| `run_scanner_heal_status_outcome_evidence.py` | dev-tool | Assembles measured Scanner/Heal G05/G06/R-D status-and-outcome release descriptors from same-run status, compatibility, and disposition artifacts | `.config/scanner-heal-required-tests.json`; `test_scanner_heal_status_outcome_evidence.sh` |
| `run_scanner_heal_maintenance_evidence.py` | dev-tool | Assembles measured Scanner/Heal G11/G13 maintenance-producer release descriptors from operator-collected proof JSON | `.config/scanner-heal-required-tests.json`; `test_scanner_heal_maintenance_evidence.sh` |
| `run_scanner_heal_w13_mrf_evidence.sh` | dev-tool | Runs the W13 durable MRF replay lanes and writes G07/G08/P4 bundle-ready evidence descriptors | `docs/testing/ci-gates.md`; `test_scanner_heal_w13_mrf_evidence.sh` |
+21 -7
View File
@@ -428,6 +428,7 @@ SCANNER_HEAL_RELEASE_SCOPED_ACK_CASES = {
"exact-generation",
"scanner-instance",
"participating-peer-set",
"cleared-count-bound",
),
"participating_peer_capability_snapshot": (
"supports-scoped-ack",
@@ -452,6 +453,7 @@ SCANNER_HEAL_RELEASE_G03_REQUIRED_TRUE_FIELDS = {
"scanner_instance_observed",
"participating_peer_set_observed",
"whole_cycle_fallback_observed",
"cleared_count_bound_observed",
),
"participating_peer_capability_snapshot": (
"capability_probe_observed",
@@ -1827,6 +1829,7 @@ def release_bundle_json_artifact_mirrored_fields(gate: str, field: str) -> tuple
"raw_entry_budget",
"max_raw_entries_per_round",
"max_objects_processed_per_round",
"bounded_work_quantum_observed",
"durable_checkpoint_committed",
"no_unbounded_tail",
))
@@ -2255,6 +2258,8 @@ def validate_release_bundle_domain_evidence(gate: str, field: str, evidence: dic
f"{gate}.{field}.max_objects_processed_per_round", 1, 4096)
require(max_raw <= budget, f"{gate}.{field} raw entries exceed fixed budget")
require(max_objects <= budget, f"{gate}.{field} processed objects exceed fixed budget")
release_bundle_bool_true(evidence.get("bounded_work_quantum_observed"),
f"{gate}.{field}.bounded_work_quantum_observed")
release_bundle_bool_true(evidence.get("durable_checkpoint_committed"),
f"{gate}.{field}.durable_checkpoint_committed")
release_bundle_bool_true(evidence.get("no_unbounded_tail"), f"{gate}.{field}.no_unbounded_tail")
@@ -2772,10 +2777,10 @@ def validate_release_bundle_artifact(bundle_path: Path, source_revision: str, ga
require(evidence.get("stale_journals_after_gc") == 0,
f"{gate}.{field} requires zero stale journals after GC")
if field == "segment_activation_preflight":
require(evidence.get("production_activation") is False,
f"{gate}.{field} must keep production activation disabled")
require(evidence.get("scanner_segment_reuse_activated") is False,
f"{gate}.{field} must prove the runtime activation gate is disabled")
require(evidence.get("production_activation") is True,
f"{gate}.{field} must prove production activation is enabled")
require(evidence.get("scanner_segment_reuse_activated") is True,
f"{gate}.{field} must prove the runtime activation gate is enabled")
evidence_exact_strings(evidence.get("proof_inputs"),
SCANNER_HEAL_SEGMENT_ACTIVATION_PROOF_INPUTS,
f"{gate}.{field}.proof_inputs")
@@ -3319,6 +3324,7 @@ def write_scanner_heal_release_bundle_fixture(root: Path, directory: Path) -> Pa
"raw_entry_budget": 8,
"max_raw_entries_per_round": 8,
"max_objects_processed_per_round": 8,
"bounded_work_quantum_observed": True,
"durable_checkpoint_committed": True,
"no_unbounded_tail": True,
})
@@ -3739,6 +3745,7 @@ class SelfTests(unittest.TestCase):
"raw_entry_budget": 8,
"max_raw_entries_per_round": 8,
"max_objects_processed_per_round": 8,
"bounded_work_quantum_observed": True,
"durable_checkpoint_committed": True,
"no_unbounded_tail": True,
})
@@ -3944,8 +3951,8 @@ class SelfTests(unittest.TestCase):
},
]
if field == "segment_activation_preflight":
evidence["production_activation"] = False
evidence["scanner_segment_reuse_activated"] = False
evidence["production_activation"] = True
evidence["scanner_segment_reuse_activated"] = True
evidence["proof_inputs"] = list(SCANNER_HEAL_SEGMENT_ACTIVATION_PROOF_INPUTS)
evidence["fail_closed_checks"] = list(SCANNER_HEAL_SEGMENT_ACTIVATION_FAIL_CLOSED_CHECKS)
if field == "cold_segment_reuse_measurement":
@@ -4447,7 +4454,14 @@ class SelfTests(unittest.TestCase):
"zero stale journals",
),
("same-window-fields", "G14", "same_window_field_evidence", lambda item: item.update({"same_window_fields": ["ec8_4_evidence", "multi_set_evidence"]}), "JSON artifact same_window_fields mismatch"),
("activation-enabled", "G11", "segment_activation_preflight", lambda item: item.update({"production_activation": True}), "production activation disabled"),
("activation-disabled", "G11", "segment_activation_preflight", lambda item: item.update({"production_activation": False}), "production activation is enabled"),
(
"activation-runtime-disabled",
"G11",
"segment_activation_preflight",
lambda item: item.update({"scanner_segment_reuse_activated": False}),
"runtime activation gate is enabled",
),
(
"activation-missing-fail-closed",
"G11",
+22 -1
View File
@@ -59,7 +59,7 @@ def expected_results(mode: str, event: str, ref: str) -> dict[str, str]:
expected.update({job: "success" if mode == "full" else "skipped" for job in CODE_JOBS})
rio = mode == "full" and event in ("schedule", "workflow_dispatch")
expected.update({job: "success" if rio else "skipped" for job in OPTIONAL_JOBS[:2]})
full = mode == "full" and (event in ("merge_group", "workflow_dispatch") or (event == "push" and ref == "refs/heads/main"))
full = mode == "full" and (event in ("merge_group", "workflow_dispatch") or (event == "push" and ref in ("refs/heads/main", "refs/heads/release")))
expected["e2e-full"] = "success" if full else "skipped"
return expected
@@ -219,6 +219,27 @@ class SelfTests(unittest.TestCase):
bad = {**good, "classify-changes": {"result": "success", "outputs": selection}}
self.assertTrue(verify_results(bad, event, "refs/heads/main"))
def test_full_e2e_gate_preserves_workflow_branch_and_event_scope(self):
for event, ref, required in (
("push", "refs/heads/main", "success"),
("push", "refs/heads/release", "success"),
("push", "refs/heads/feature", "skipped"),
("push", "refs/heads/release-candidate", "skipped"),
("push", "refs/tags/release", "skipped"),
("pull_request", "refs/pull/1/merge", "skipped"),
("schedule", "refs/heads/release", "skipped"),
("workflow_dispatch", "refs/heads/feature", "success"),
("merge_group", "refs/heads/gh-readonly-queue/release/pr-1", "success"),
):
with self.subTest(event=event, ref=ref):
expected = expected_results("full", event, ref)
self.assertEqual(expected["e2e-full"], required)
needs = {job: {"result": result} for job, result in expected.items()}
needs["classify-changes"]["outputs"] = {"mode": "full"}
for result in ("success", "skipped", "failure", "cancelled"):
needs["e2e-full"]["result"] = result
self.assertEqual(verify_results(needs, event, ref) == [], result == required)
def test_repository_wiring_and_missing_dependency_regression(self):
self.assertEqual(check_workflow(ROOT), [])
with tempfile.TemporaryDirectory() as directory:
@@ -74,10 +74,31 @@ def converged(report, objects):
def replays_raw_window(previous, current):
if (previous["raw_entries"] == 0 or current["raw_entries"] == 0
or previous["raw_first_entry"] is None or current["raw_first_entry"] is None
or previous["raw_last_entry"] is None or current["raw_last_entry"] is None):
return False
raw_page_index_progressed = (
previous["raw_page_index_parent"] == current["raw_page_index_parent"]
and (
current["raw_page_index_committed_entries"] > previous["raw_page_index_committed_entries"]
or (current["raw_page_index_complete"] and not previous["raw_page_index_complete"])
)
)
return (previous["raw_first_entry"] == current["raw_first_entry"]
and previous["raw_last_entry"] == current["raw_last_entry"]
and previous["objects_retained"] == current["objects_before"]
and current["objects_retained"] == previous["objects_retained"])
and current["objects_retained"] == previous["objects_retained"]
and not raw_page_index_progressed)
def fully_retained(report, objects):
return all(report[key] == objects for key in
("objects_retained", "versions_retained", "bytes_retained"))
def final_complete_recheck_after_full_retention(previous, current, objects):
return fully_retained(previous, objects) and converged(current, objects)
def validate_recoverable_quantum(reports, *, objects, budget, require_converged):
@@ -97,7 +118,8 @@ def validate_recoverable_quantum(reports, *, objects, budget, require_converged)
raise ValueError("durable retained coverage did not survive process restart")
if report["objects_retained"] < previous["objects_retained"]:
raise ValueError("durable retained coverage regressed across restart")
if replays_raw_window(previous, report):
if (not final_complete_recheck_after_full_retention(previous, report, objects)
and replays_raw_window(previous, report)):
raise ValueError("raw enumeration window replayed without durable coverage")
if (report["raw_page_index_parent"] == previous["raw_page_index_parent"]
and report["raw_page_index_committed_entries"] < previous["raw_page_index_committed_entries"]
@@ -32,6 +32,11 @@ def positive_int(value: Any, name: str, minimum: int = 1) -> int:
return value
def reject_non_measured_markers(payload: dict[str, Any], label: str) -> None:
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
require(payload.get(marker) is not True, f"{label} is {marker}")
def timestamp(value: Any, name: str) -> str:
require(isinstance(value, str) and value.endswith("Z"), f"invalid {name}")
datetime.fromisoformat(value.replace("Z", "+00:00"))
@@ -63,6 +68,7 @@ def load_reports(directory: Path) -> list[dict[str, Any]]:
"raw_page_index_indexed_entries",
}
for index, report in enumerate(reports):
reject_non_measured_markers(report, f"round {index}")
require(report.get("schema") == 1, f"round {index} has wrong schema")
require(report.get("round") == index, f"round {index} order mismatch")
for key in (
@@ -88,8 +94,7 @@ def load_reports(directory: Path) -> list[dict[str, Any]]:
def require_measured_manifest(path: Path, source_revision: str) -> dict[str, Any]:
manifest = read_json(path)
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
require(manifest.get(marker) is not True, f"checkpoint/crash manifest is {marker}")
reject_non_measured_markers(manifest, "checkpoint/crash manifest")
require(manifest.get("schema") == 1, "unsupported manifest schema")
require(manifest.get("evidence_type") == "measured", "manifest must be measured")
require(manifest.get("source_revision") == source_revision, "manifest source revision mismatch")
@@ -108,8 +113,7 @@ def derived_measured_manifest(directory: Path, source_revision: str, reports: li
request_path = directory / "request.json"
request = read_json(request_path)
require(isinstance(request, dict), "diagnostic request must be a JSON object")
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
require(request.get(marker) is not True, f"diagnostic request is {marker}")
reject_non_measured_markers(request, "diagnostic request")
objects = positive_int(request.get("objects"), "request.objects")
raw_entry_budget = positive_int(request.get("raw_entry_budget"), "request.raw_entry_budget")
final_round = positive_int(request.get("round"), "request.round", 0)
@@ -251,6 +255,7 @@ def build_descriptor(args: argparse.Namespace) -> Path:
"raw_entry_budget": summary["raw_entry_budget"],
"max_raw_entries_per_round": summary["max_raw_entries_per_round"],
"max_objects_processed_per_round": summary["max_objects_processed_per_round"],
"bounded_work_quantum_observed": True,
"durable_checkpoint_committed": True,
"no_unbounded_tail": True,
}),
@@ -422,6 +427,24 @@ def run_self_test() -> None:
else:
raise ValueError("self-test accepted non-converged diagnostic")
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
source_revision = git_head()
manifest, diagnostic = write_self_test_inputs(root, source_revision)
report = read_json(diagnostic / "round-0.json")
report["synthetic"] = True
write_json(diagnostic / "round-0.json", report)
try:
build_descriptor(parse_args([
"--manifest", str(manifest),
"--diagnostic-dir", str(diagnostic),
"--out-dir", str(root / "out"),
]))
except ValueError as err:
require("round 0 is synthetic" in str(err), "wrong self-test failure for synthetic round report")
else:
raise ValueError("self-test accepted synthetic diagnostic round report")
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
@@ -48,6 +48,11 @@ def positive_int(value: Any, name: str, minimum: int = 1) -> int:
return value
def reject_non_measured_markers(payload: dict[str, Any], label: str) -> None:
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
require(payload.get(marker) is not True, f"{label} is {marker}")
def parse_case_dir_arg(value: str) -> tuple[str | None, Path]:
if "=" in value:
case_id, raw_path = value.split("=", 1)
@@ -72,8 +77,7 @@ def proof_case_artifact_path(proof_path: Path, sample: dict[str, Any], index: in
def load_proof(path: Path, source_revision: str) -> dict[str, Any]:
proof = read_json(path)
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
require(proof.get(marker) is not True, f"G14 proof is {marker}")
reject_non_measured_markers(proof, "G14 proof")
require(proof.get("schema") == 1, "unsupported G14 proof schema")
require(proof.get("evidence_type") == "measured", "G14 proof must be measured")
require(proof.get("source_revision") == source_revision, "G14 proof source revision mismatch")
@@ -110,8 +114,7 @@ def load_proof(path: Path, source_revision: str) -> dict[str, Any]:
f"G14 case evidence {index} measurement window mismatch")
artifact_payload = read_json(artifact)
require(isinstance(artifact_payload, dict), f"G14 case evidence {index} artifact must be a JSON object")
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
require(artifact_payload.get(marker) is not True, f"G14 case evidence {index} artifact is {marker}")
reject_non_measured_markers(artifact_payload, f"G14 case evidence {index} artifact")
require(artifact_payload.get("case") == sample["case"], f"G14 case evidence {index} artifact case mismatch")
if "source_revision" in artifact_payload:
require(artifact_payload["source_revision"] == source_revision,
@@ -127,6 +130,10 @@ def load_case_directory(raw_value: str, source_revision: str) -> dict[str, Any]:
require(directory.is_dir(), f"G14 case directory is missing: {directory}")
run = read_json(directory / "run.json")
execution = read_json(directory / "execution.json")
require(isinstance(run, dict), "G14 case run must be a JSON object")
require(isinstance(execution, dict), "G14 case execution must be a JSON object")
reject_non_measured_markers(run, "G14 case run")
reject_non_measured_markers(execution, "G14 case execution")
require(run.get("schema") == 1, "G14 case run schema mismatch")
require(isinstance(run.get("run_id"), str) and re.fullmatch(r"[0-9a-f]{32}", run["run_id"]),
"invalid G14 case run id")
@@ -156,6 +163,7 @@ def load_case_directory(raw_value: str, source_revision: str) -> dict[str, Any]:
if expected_case is not None:
require(oracle.get("case") == expected_case, "G14 case directory case id mismatch")
require(oracle.get("source_revision") == source_revision, "G14 oracle source revision mismatch")
reject_non_measured_markers(oracle, "G14 oracle")
require(oracle.get("evidence") in {"process-restart", "process-crash-restart"}, "G14 oracle has wrong evidence type")
topology = oracle.get("topology")
require(isinstance(topology, dict), "G14 oracle missing topology")
@@ -504,6 +512,27 @@ def run_self_test() -> None:
else:
raise ValueError("self-test accepted case evidence without multi-pool proof")
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
source_revision = git_head()
multi_pool = write_self_test_case_dir(root, source_revision, "background-target-crash-ec8-4-multi-pool", 3, 3)
oracle_path = multi_pool / "background-target-crash-ec8-4-multi-pool.json"
oracle = read_json(oracle_path)
oracle["fixture_only"] = True
write_json(oracle_path, oracle)
execution = read_json(multi_pool / "execution.json")
execution["artifacts"][oracle_path.name] = digest(oracle_path)
write_json(multi_pool / "execution.json", execution)
try:
build_descriptor(parse_args([
"--case-dir", f"background-target-crash-ec8-4-multi-pool={multi_pool}",
"--out-dir", str(root / "out"),
]))
except ValueError as err:
require("G14 oracle is fixture_only" in str(err), "wrong self-test failure for fixture oracle")
else:
raise ValueError("self-test accepted fixture-only G14 case oracle")
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
@@ -146,13 +146,16 @@ def build_descriptor(args: argparse.Namespace) -> Path:
"gates": gates,
})
for gate in ("G11", "G13"):
subprocess.check_call([
check = subprocess.run([
sys.executable,
str(ROOT / "scripts/check_test_wiring.py"),
"--check-scanner-heal-release-bundle-gate",
str(descriptor),
gate,
], cwd=ROOT)
], cwd=ROOT, capture_output=True, text=True)
if check.returncode != 0:
details = "\n".join(part for part in (check.stdout.strip(), check.stderr.strip()) if part)
raise ValueError(details or f"{gate} release bundle gate check failed")
return descriptor
@@ -188,8 +191,8 @@ def write_self_test_proof(path: Path, source_revision: str) -> None:
"unknown_producer_excluded": True,
},
"segment_activation_preflight": {
"production_activation": False,
"scanner_segment_reuse_activated": False,
"production_activation": True,
"scanner_segment_reuse_activated": True,
"proof_inputs": list(wiring.SCANNER_HEAL_SEGMENT_ACTIVATION_PROOF_INPUTS),
"fail_closed_checks": list(wiring.SCANNER_HEAL_SEGMENT_ACTIVATION_FAIL_CLOSED_CHECKS),
},
@@ -243,6 +246,21 @@ def run_self_test() -> None:
else:
raise ValueError("self-test accepted synthetic proof")
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
source_revision = git_head()
proof = root / "maintenance-proof.json"
write_self_test_proof(proof, source_revision)
payload = wiring.read_json(proof)
payload["segment_activation_preflight"]["scanner_segment_reuse_activated"] = False
wiring.write_json(proof, payload)
try:
build_descriptor(parse_args(["--proof-json", str(proof), "--out-dir", str(root / "out")]))
except ValueError as err:
wiring.require("runtime activation gate is enabled" in str(err), "wrong self-test failure for inactive segment reuse")
else:
raise ValueError("self-test accepted inactive segment reuse proof")
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
+470
View File
@@ -0,0 +1,470 @@
#!/usr/bin/env python3
"""Collect Scanner/Heal status-and-outcome raw evidence from live observations.
This helper normalizes operator-collected live observation JSON into the three
measured raw artifacts consumed by run_scanner_heal_status_outcome_evidence.py.
It does not approve a release bundle by itself.
"""
from __future__ import annotations
import argparse
from datetime import datetime, timezone
import json
import re
import subprocess
import sys
from pathlib import Path
from typing import Any
import check_test_wiring as wiring
ROOT = Path(__file__).resolve().parents[1]
def git_head() -> str:
return subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=ROOT, text=True).strip()
def timestamp(value: Any, name: str) -> str:
wiring.require(isinstance(value, str) and value.endswith("Z"), f"invalid {name}")
parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
wiring.require(parsed.tzinfo is not None, f"{name} must include timezone")
return parsed.isoformat().replace("+00:00", "Z")
def identity_string(value: Any, name: str) -> str:
wiring.require(
isinstance(value, str) and re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._:-]{7,127}", value),
f"invalid {name}",
)
return value
def measured_observation(path: Path, source_revision: str) -> dict[str, Any]:
payload = wiring.read_json(path.resolve())
wiring.require(isinstance(payload, dict), "observation must be a JSON object")
for marker in ("fixture", "fixture_only", "dry_run", "synthetic"):
wiring.require(payload.get(marker) is not True, f"observation is {marker}")
wiring.require(payload.get("schema") == 1, "observation schema must be 1")
wiring.require(payload.get("evidence_type") == "measured", "observation must be measured")
wiring.require(payload.get("source_revision") == source_revision, "observation source revision mismatch")
run_id = identity_string(payload.get("run_id"), "run_id")
window_id = identity_string(payload.get("measurement_window_id"), "measurement_window_id")
wiring.require(run_id != window_id, "run/window identities must differ")
started = timestamp(payload.get("started_at"), "started_at")
finished = timestamp(payload.get("finished_at"), "finished_at")
wiring.require(
datetime.fromisoformat(finished.replace("Z", "+00:00"))
>= datetime.fromisoformat(started.replace("Z", "+00:00")),
"finished_at precedes started_at",
)
wiring.require(isinstance(payload.get("command"), list) and payload["command"], "missing command provenance")
return payload
def object_items(payload: dict[str, Any], key: str) -> list[dict[str, Any]]:
items = payload.get(key)
wiring.require(isinstance(items, list) and items, f"missing {key}")
for index, item in enumerate(items):
wiring.require(isinstance(item, dict), f"{key}[{index}] must be an object")
return items
def case_map(items: list[dict[str, Any]], key: str, expected: tuple[str, ...]) -> dict[str, dict[str, Any]]:
observed: dict[str, dict[str, Any]] = {}
for item in items:
case = item.get("case")
wiring.require(isinstance(case, str) and case, f"{key} item missing case")
wiring.require(case not in observed, f"{key} duplicate case: {case}")
observed[case] = item
missing = [case for case in expected if case not in observed]
unknown = [case for case in observed if case not in expected]
wiring.require(not missing, f"{key} missing cases: {', '.join(missing)}")
wiring.require(not unknown, f"{key} unknown cases: {', '.join(unknown)}")
return observed
def bool_true(value: Any, name: str) -> None:
wiring.release_bundle_bool_true(value, name)
def status_outcome(payload: dict[str, Any], common: dict[str, Any]) -> dict[str, Any]:
outcome_items = case_map(
object_items(payload, "status_outcomes"),
"status_outcomes",
wiring.SCANNER_HEAL_RELEASE_G05_PER_OBJECT_OUTCOME_CASES,
)
counts = {"repaired": 0, "healthy": 0, "skipped": 0, "failed": 0}
for case, item in outcome_items.items():
outcome = item.get("outcome")
wiring.require(outcome in counts, f"status_outcomes.{case} has invalid outcome")
bool_true(item.get("status_matches_object_oracle"), f"status_outcomes.{case}.status_matches_object_oracle")
counts[outcome] += 1
for outcome, count in counts.items():
wiring.require(count > 0, f"missing measured {outcome} outcome")
retention = case_map(
object_items(payload, "terminal_retention_samples"),
"terminal_retention_samples",
wiring.SCANNER_HEAL_RELEASE_G05_TERMINAL_RETENTION_CASES,
)
window = wiring.evidence_integer(
payload.get("terminal_retention_window_seconds"),
"terminal_retention_window_seconds",
1,
86400,
)
max_age = 0
pruned = 0
for case, item in retention.items():
bool_true(item.get("retained_until_window"), f"terminal_retention_samples.{case}.retained_until_window")
max_age = max(
max_age,
wiring.evidence_integer(
item.get("max_age_seconds"),
f"terminal_retention_samples.{case}.max_age_seconds",
0,
86400,
),
)
if item.get("pruned_after_window") is True:
pruned += 1
wiring.require(max_age <= window, "terminal retention max age exceeds window")
wiring.require(pruned > 0, "no terminal records were pruned after the retention window")
return {
**common,
"per_object_outcome_cases": list(wiring.SCANNER_HEAL_RELEASE_G05_PER_OBJECT_OUTCOME_CASES),
"outcome_counts": counts,
"status_matches_object_oracle": True,
"terminal_retention_cases": list(wiring.SCANNER_HEAL_RELEASE_G05_TERMINAL_RETENTION_CASES),
"terminal_retention_window_seconds": window,
"max_terminal_record_age_seconds": max_age,
"terminal_records_pruned_after_window": pruned,
}
def status_compat(payload: dict[str, Any], common: dict[str, Any]) -> dict[str, Any]:
status = case_map(
object_items(payload, "status_samples"),
"status_samples",
wiring.SCANNER_HEAL_RELEASE_G06_CONCURRENT_STATUS_CASES,
)
status_samples = 0
degraded = False
for case, item in status.items():
bool_true(item.get("http_success"), f"status_samples.{case}.http_success")
status_samples += wiring.evidence_integer(item.get("samples"), f"status_samples.{case}.samples", 1, 2**31 - 1)
degraded = degraded or item.get("degraded") is True
legacy = case_map(
object_items(payload, "legacy_client_samples"),
"legacy_client_samples",
wiring.SCANNER_HEAL_RELEASE_G06_LEGACY_CLIENT_CASES,
)
for case, item in legacy.items():
bool_true(item.get("accepted"), f"legacy_client_samples.{case}.accepted")
rustfs_and_minio = (
legacy["rustfs-admin-v3-background-heal-status"].get("path_compatible") is True
and legacy["minio-admin-v3-background-heal-status"].get("path_compatible") is True
)
empty_body = legacy["heal-client-token-empty-body"].get("empty_body_accepted") is True
bool_true(rustfs_and_minio, "rustfs and minio status path compatibility")
bool_true(empty_body, "empty body heal status request")
truncation = case_map(
object_items(payload, "truncation_samples"),
"truncation_samples",
wiring.SCANNER_HEAL_RELEASE_G06_TRUNCATION_CASES,
)
max_payload = 1
for case, item in truncation.items():
bool_true(item.get("rejected"), f"truncation_samples.{case}.rejected")
max_payload = max(
max_payload,
wiring.evidence_integer(item.get("payload_bytes"), f"truncation_samples.{case}.payload_bytes", 1, 2**20),
)
return {
**common,
"concurrent_status_cases": list(wiring.SCANNER_HEAL_RELEASE_G06_CONCURRENT_STATUS_CASES),
"status_samples": status_samples,
"all_status_responses_http_success": True,
"partial_status_reports_degraded": degraded,
"legacy_client_cases": list(wiring.SCANNER_HEAL_RELEASE_G06_LEGACY_CLIENT_CASES),
"rustfs_and_minio_paths_compatible": rustfs_and_minio,
"empty_body_status_requests_accepted": empty_body,
"truncation_cases": list(wiring.SCANNER_HEAL_RELEASE_G06_TRUNCATION_CASES),
"truncated_payloads_rejected": True,
"max_status_payload_bytes": max_payload,
}
def disposition(payload: dict[str, Any], common: dict[str, Any]) -> dict[str, Any]:
managers = case_map(
object_items(payload, "manager_dispositions"),
"manager_dispositions",
wiring.SCANNER_HEAL_RELEASE_RD_MANAGER_CASES,
)
for case, item in managers.items():
bool_true(item.get("terminal"), f"manager_dispositions.{case}.terminal")
events = case_map(
object_items(payload, "event_dispositions"),
"event_dispositions",
wiring.SCANNER_HEAL_RELEASE_RD_EVENT_CASES,
)
for case, item in events.items():
bool_true(item.get("correlates_to_manager"), f"event_dispositions.{case}.correlates_to_manager")
ledgers = case_map(
object_items(payload, "ledger_dispositions"),
"ledger_dispositions",
wiring.SCANNER_HEAL_RELEASE_RD_LEDGER_CASES,
)
for case, item in ledgers.items():
bool_true(item.get("correlates_to_events"), f"ledger_dispositions.{case}.correlates_to_events")
bool_true(item.get("replay_preserves_terminal"), f"ledger_dispositions.{case}.replay_preserves_terminal")
grace = case_map(
object_items(payload, "grace_samples"),
"grace_samples",
wiring.SCANNER_HEAL_RELEASE_RD_GRACE_CASES,
)
grace_window = wiring.evidence_integer(payload.get("grace_window_seconds"), "grace_window_seconds", 1, 86400)
retained = False
pruned = False
for item in grace.values():
retained = retained or item.get("retention_observed") is True
pruned = pruned or item.get("expiry_pruned_terminal_records") is True
bool_true(retained, "grace retention observed")
bool_true(pruned, "grace expiry pruned terminal records")
return {
**common,
"manager_disposition_cases": list(wiring.SCANNER_HEAL_RELEASE_RD_MANAGER_CASES),
"manager_dispositions_are_terminal": True,
"event_disposition_cases": list(wiring.SCANNER_HEAL_RELEASE_RD_EVENT_CASES),
"events_correlate_to_manager_dispositions": True,
"ledger_disposition_cases": list(wiring.SCANNER_HEAL_RELEASE_RD_LEDGER_CASES),
"ledger_correlates_to_events": True,
"ledger_replay_preserves_terminal_disposition": True,
"grace_cases": list(wiring.SCANNER_HEAL_RELEASE_RD_GRACE_CASES),
"grace_window_seconds": grace_window,
"grace_retention_observed": retained,
"grace_expiry_pruned_terminal_records": pruned,
}
def common_raw(payload: dict[str, Any]) -> dict[str, Any]:
return {
"schema": 1,
"evidence_type": "measured",
"source_revision": payload["source_revision"],
"run_id": payload["run_id"],
"measurement_window_id": payload["measurement_window_id"],
"started_at": payload["started_at"],
"finished_at": payload["finished_at"],
"command": payload["command"],
}
def collect(args: argparse.Namespace) -> tuple[Path, Path, Path]:
out_dir = args.out_dir.resolve()
wiring.require(not out_dir.exists(), "output directory must be new")
source_revision = args.source_revision or git_head()
observed = measured_observation(args.observations_json, source_revision)
common = common_raw(observed)
out_dir.mkdir(parents=True)
status_outcome_path = out_dir / "status-outcome.json"
status_compat_path = out_dir / "status-compat.json"
disposition_path = out_dir / "disposition.json"
wiring.write_json(status_outcome_path, status_outcome(observed, common))
wiring.write_json(status_compat_path, status_compat(observed, common))
wiring.write_json(disposition_path, disposition(observed, common))
return status_outcome_path, status_compat_path, disposition_path
def write_self_test_observation(path: Path, source_revision: str) -> None:
now = datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
payload = {
"schema": 1,
"evidence_type": "measured",
"source_revision": source_revision,
"run_id": f"status-probe-{source_revision[:12]}",
"measurement_window_id": f"status-probe-window-{source_revision[:12]}",
"started_at": now,
"finished_at": now,
"command": ["scripts/run_scanner_heal_status_outcome_probe.py", "--observations-json", "<live-observations>"],
"terminal_retention_window_seconds": 3600,
"grace_window_seconds": 300,
"status_outcomes": [
{"case": "object-repaired", "outcome": "repaired", "status_matches_object_oracle": True},
{"case": "object-already-healthy", "outcome": "healthy", "status_matches_object_oracle": True},
{"case": "object-skipped-by-policy", "outcome": "skipped", "status_matches_object_oracle": True},
{"case": "object-failed-and-retained", "outcome": "failed", "status_matches_object_oracle": True},
],
"terminal_retention_samples": [
{"case": "finished-retained-until-window", "retained_until_window": True, "max_age_seconds": 1200},
{"case": "failed-retained-until-window", "retained_until_window": True, "max_age_seconds": 1300},
{"case": "canceled-retained-until-window", "retained_until_window": True, "max_age_seconds": 1400},
{
"case": "expired-terminal-pruned-after-window",
"retained_until_window": True,
"max_age_seconds": 3599,
"pruned_after_window": True,
},
],
"status_samples": [
{"case": "status-during-admin-heal", "samples": 2, "http_success": True},
{"case": "status-during-background-heal", "samples": 2, "http_success": True},
{"case": "status-while-peer-down", "samples": 2, "http_success": True, "degraded": True},
{"case": "status-after-peer-rejoin", "samples": 2, "http_success": True},
],
"legacy_client_samples": [
{"case": "rustfs-admin-v3-background-heal-status", "accepted": True, "path_compatible": True},
{"case": "minio-admin-v3-background-heal-status", "accepted": True, "path_compatible": True},
{"case": "heal-client-token-empty-body", "accepted": True, "empty_body_accepted": True},
{"case": "node-heal-status-v1-wire", "accepted": True},
],
"truncation_samples": [
{"case": "oversize-node-status-reject", "rejected": True, "payload_bytes": 1048576},
{"case": "truncated-node-status-reject", "rejected": True, "payload_bytes": 4096},
{"case": "trailing-data-node-status-reject", "rejected": True, "payload_bytes": 4096},
],
"manager_dispositions": [
{"case": "accepted", "terminal": True},
{"case": "coalesced-duplicate", "terminal": True},
{"case": "rejected-policy", "terminal": True},
{"case": "terminal-retained", "terminal": True},
],
"event_dispositions": [
{"case": "event-repaired", "correlates_to_manager": True},
{"case": "event-failed", "correlates_to_manager": True},
{"case": "event-skipped", "correlates_to_manager": True},
{"case": "event-grace-retained", "correlates_to_manager": True},
],
"ledger_dispositions": [
{"case": "ledger-recorded", "correlates_to_events": True, "replay_preserves_terminal": True},
{"case": "ledger-replayed", "correlates_to_events": True, "replay_preserves_terminal": True},
{"case": "ledger-discharged", "correlates_to_events": True, "replay_preserves_terminal": True},
{"case": "ledger-pruned-after-grace", "correlates_to_events": True, "replay_preserves_terminal": True},
],
"grace_samples": [
{"case": "grace-open-retains-disposition", "retention_observed": True},
{"case": "grace-expired-prunes-terminal", "expiry_pruned_terminal_records": True},
{"case": "restart-preserves-grace-clock", "retention_observed": True},
],
}
wiring.write_json(path, payload)
def run_self_test() -> None:
import tempfile
source_revision = git_head()
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
observation = root / "observation.json"
write_self_test_observation(observation, source_revision)
raw_paths = collect(parse_args([
"--observations-json",
str(observation),
"--out-dir",
str(root / "raw"),
]))
descriptor = root / "descriptor"
subprocess.check_call([
sys.executable,
str(ROOT / "scripts/run_scanner_heal_status_outcome_evidence.py"),
"--status-outcome-json",
str(raw_paths[0]),
"--status-compat-json",
str(raw_paths[1]),
"--disposition-json",
str(raw_paths[2]),
"--out-dir",
str(descriptor),
], cwd=ROOT)
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
observation = root / "observation.json"
write_self_test_observation(observation, source_revision)
payload = wiring.read_json(observation)
payload["status_outcomes"].pop()
wiring.write_json(observation, payload)
try:
collect(parse_args(["--observations-json", str(observation), "--out-dir", str(root / "raw")]))
except ValueError as err:
wiring.require("status_outcomes missing cases" in str(err), "wrong self-test failure for missing outcome")
else:
raise ValueError("self-test accepted incomplete status outcome observations")
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
observation = root / "observation.json"
write_self_test_observation(observation, source_revision)
payload = wiring.read_json(observation)
payload["synthetic"] = True
wiring.write_json(observation, payload)
try:
collect(parse_args(["--observations-json", str(observation), "--out-dir", str(root / "raw")]))
except ValueError as err:
wiring.require("observation is synthetic" in str(err), "wrong self-test failure for synthetic observation")
else:
raise ValueError("self-test accepted synthetic status/outcome observations")
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
observation = root / "observation.json"
write_self_test_observation(observation, source_revision)
payload = wiring.read_json(observation)
payload["measurement_window_id"] = payload["run_id"]
wiring.write_json(observation, payload)
try:
collect(parse_args(["--observations-json", str(observation), "--out-dir", str(root / "raw")]))
except ValueError as err:
wiring.require("run/window identities must differ" in str(err), "wrong self-test failure for identity reuse")
else:
raise ValueError("self-test accepted reused run/window identities")
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--observations-json", type=Path)
parser.add_argument("--out-dir", type=Path)
parser.add_argument("--source-revision")
parser.add_argument("--self-test", action="store_true")
args = parser.parse_args(argv)
if not args.self_test:
if args.observations_json is None:
parser.error("--observations-json is required unless --self-test is used")
if args.out_dir is None:
parser.error("--out-dir is required unless --self-test is used")
return args
def main() -> None:
args = parse_args()
if args.self_test:
run_self_test()
return
paths = collect(args)
json.dump(
{
"status_outcome_json": str(paths[0]),
"status_compat_json": str(paths[1]),
"disposition_json": str(paths[2]),
},
sys.stdout,
indent=2,
allow_nan=False,
)
sys.stdout.write("\n")
if __name__ == "__main__":
main()
@@ -4,6 +4,8 @@ import unittest
from diagnose_scanner_enumeration_restart import (
converged,
final_complete_recheck_after_full_retention,
fully_retained,
replays_raw_window,
validate_recoverable_quantum,
validate_report,
@@ -29,6 +31,7 @@ class ReportTests(unittest.TestCase):
report = self.report()
self.validate(report)
self.assertTrue(converged(report, 4))
self.assertTrue(fully_retained(report, 4))
def test_incomplete_or_inexact_coverage_cannot_pass(self):
for key, value in (("snapshot_complete", False), ("objects_retained", 3),
@@ -132,6 +135,9 @@ class ReportTests(unittest.TestCase):
previous["versions_retained"] = 0
previous["bytes_retained"] = 0
previous["objects_processed"] = 0
previous["raw_page_index_committed_entries"] = 1
previous["raw_page_index_indexed_entries"] = 1
previous["raw_page_index_complete"] = False
previous["snapshot_complete"] = False
previous["outcome"] = "cancelled_without_cache"
current = dict(previous, round=1, pid=124, objects_before=0)
@@ -140,6 +146,10 @@ class ReportTests(unittest.TestCase):
advanced = dict(current, objects_retained=1)
self.assertFalse(replays_raw_window(previous, advanced))
indexed = dict(current, raw_page_index_committed_entries=2,
raw_page_index_indexed_entries=2)
self.assertFalse(replays_raw_window(previous, indexed))
def test_recoverable_quantum_rejects_replayed_raw_window(self):
previous = self.report()
previous.update(objects_retained=0, versions_retained=0, bytes_retained=0,
@@ -149,6 +159,35 @@ class ReportTests(unittest.TestCase):
with self.assertRaisesRegex(ValueError, "raw enumeration window replayed"):
validate_recoverable_quantum([previous, current], objects=4, budget=16, require_converged=False)
def test_recoverable_quantum_allows_final_complete_round_after_full_retention(self):
previous = self.report()
previous.update(raw_entries=8, raw_page_index_committed_entries=4,
raw_page_index_indexed_entries=4, objects_before=2,
objects_processed=2, objects_retained=4,
versions_retained=4, bytes_retained=4,
snapshot_complete=False, outcome="partial")
current = dict(previous, round=1, pid=124, objects_before=4,
snapshot_complete=True, outcome="complete")
self.assertTrue(replays_raw_window(previous, current))
self.assertTrue(final_complete_recheck_after_full_retention(previous, current, 4))
validate_recoverable_quantum([previous, current], objects=4, budget=16, require_converged=True)
def test_recoverable_quantum_rejects_partial_replay_after_full_retention(self):
previous = self.report()
previous.update(raw_entries=8, raw_page_index_committed_entries=4,
raw_page_index_indexed_entries=4, objects_before=2,
objects_processed=2, objects_retained=4,
versions_retained=4, bytes_retained=4,
snapshot_complete=False, outcome="partial")
current = dict(previous, round=1, pid=124, objects_before=4,
snapshot_complete=False, outcome="partial")
self.assertTrue(replays_raw_window(previous, current))
self.assertFalse(final_complete_recheck_after_full_retention(previous, current, 4))
with self.assertRaisesRegex(ValueError, "raw enumeration window replayed"):
validate_recoverable_quantum([previous, current], objects=4, budget=16, require_converged=False)
def test_recoverable_quantum_requires_three_stage_progress_and_convergence(self):
first = self.report()
first.update(round=0, pid=123, raw_entries=2, raw_page_index_committed_entries=2,
@@ -3,5 +3,7 @@ set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
RUNNER="$SCRIPT_DIR/run_scanner_heal_status_outcome_evidence.py"
PROBE="$SCRIPT_DIR/run_scanner_heal_status_outcome_probe.py"
"${RUSTFS_PYTHON_BIN:-python3}" "$PROBE" --self-test
"${RUSTFS_PYTHON_BIN:-python3}" "$RUNNER" --self-test