fix: address release acceptance regressions (#7541)

This commit is contained in:
Zhengchao An
2026-09-09 08:33:59 +08:00
committed by GitHub
parent 8617f2701b
commit e24eae9eaa
29 changed files with 2582 additions and 134 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=f0c78fdb93471575d9a64c5c46eae6c806bdd0bc10a6e33d7fb574aabd8db5a3
sha256-linux=03ed7016cab672de9320e31375a0358eceacb4408b0e79cf063614fa7c878b87
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
sha256-linux=9515861be899ceb10e2e0ef93c34208bb7a7a8a7f8067a02db4cfba23270ebd6
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=364f2329a7b72eb9f1608dbe1a3af37af4095354014f3cbe23ca448492d89961
sha256-linux=60983f1ebe7068cf660d473c5f76c76a650410ccc99d71934ddca7fd67607987
sha256-darwin=83a7dcaffd5a789517ae9f02a224f66a9713937885cff96fca2ad7e216f197ae
sha256-linux=626c10f8c964507ff987b6c86069e9019dc6d2ae7fb02db9be5df5aa8cc5145b
@@ -268,6 +268,7 @@ async fn test_bucket_cors_write_is_visible_on_peer_before_response() -> Result<(
let rule = CorsRule::builder()
.allowed_methods("GET")
.allowed_origins("https://example.com")
.allowed_headers("*")
.build()?;
let configuration = CorsConfiguration::builder().cors_rules(rule).build()?;
@@ -288,6 +289,60 @@ async fn test_bucket_cors_write_is_visible_on_peer_before_response() -> Result<(
assert_eq!(rules[0].allowed_methods(), ["GET"]);
assert_eq!(rules[0].allowed_origins(), ["https://example.com"]);
let http = reqwest::Client::builder().no_proxy().build()?;
let url = format!("http://{}/{}", cluster.nodes[1].address, BUCKET_METADATA_RELOAD_BUCKET);
let without_headers = http
.request(reqwest::Method::OPTIONS, &url)
.header("Origin", "https://example.com")
.header("Access-Control-Request-Method", "GET")
.send()
.await?;
assert!(without_headers.status().is_success());
assert!(!without_headers.headers().contains_key("access-control-allow-headers"));
assert!(
without_headers
.headers()
.get("vary")
.and_then(|value| value.to_str().ok())
.is_some_and(|value| value.contains("Access-Control-Request-Headers")),
"a cached header-free preflight must not suppress a later requested header grant"
);
let preflight = http
.request(reqwest::Method::OPTIONS, &url)
.header("Origin", "https://example.com")
.header("Access-Control-Request-Method", "GET")
.header("Access-Control-Request-Headers", "X-Another-Header, x-could-be-anything")
.send()
.await?;
assert!(preflight.status().is_success(), "peer preflight should succeed: {preflight:?}");
assert_eq!(
preflight
.headers()
.get("access-control-allow-headers")
.and_then(|value| value.to_str().ok()),
Some("x-another-header,x-could-be-anything"),
"a wildcard rule must return only the headers requested by this preflight"
);
assert!(
preflight
.headers()
.get("vary")
.and_then(|value| value.to_str().ok())
.is_some_and(|value| value.contains("Access-Control-Request-Headers")),
"preflight caches must distinguish the requested header list"
);
let denied = http
.request(reqwest::Method::OPTIONS, &url)
.header("Origin", "https://disallowed.example.com")
.header("Access-Control-Request-Method", "GET")
.header("Access-Control-Request-Headers", "x-another-header")
.send()
.await?;
assert!(
!denied.headers().contains_key("access-control-allow-headers"),
"a rejected origin must not receive the requested header grant"
);
writer
.delete_bucket_cors()
.bucket(BUCKET_METADATA_RELOAD_BUCKET)
@@ -962,27 +962,6 @@ pub(crate) async fn wait_for_rebalance_active(
}
}
pub(crate) async fn wait_for_rebalance_running_with_progress(
cluster: &RustFSTestClusterEnvironment,
expected_id: &str,
timeout: Duration,
) -> TestResult {
let deadline = Instant::now() + timeout;
loop {
let status = rebalance_status_json(cluster).await?;
if rebalance_running_with_progress(&status, expected_id)? {
return Ok(());
}
if Instant::now() >= deadline {
return Err(format!(
"rebalance did not become active with non-zero progress within {timeout:?}; last status: {status}"
)
.into());
}
sleep(Duration::from_millis(100)).await;
}
}
pub(crate) async fn wait_for_rebalance_complete(
cluster: &RustFSTestClusterEnvironment,
expected_id: &str,
@@ -14,9 +14,9 @@
use super::harness::{
DECOMMISSION_POOL_ID, DistCluster, DistLayout, TestResult, assert_inventory, decommission_running_with_progress,
decommission_status_json, put_inventory_retrying, rebalance_running_with_progress, rebalance_status_json,
retrying_get_equals, retrying_put, start_decommission, start_rebalance, unique_bucket, wait_for_decommission_complete,
wait_for_decommission_running_with_progress, wait_for_rebalance_complete, wait_for_rebalance_running_with_progress,
decommission_status_json, put_inventory_retrying, rebalance_active, rebalance_status_json, retrying_get_equals, retrying_put,
start_decommission, start_rebalance, unique_bucket, wait_for_decommission_complete,
wait_for_decommission_running_with_progress, wait_for_rebalance_active, wait_for_rebalance_complete,
};
use crate::common::init_logging;
use std::time::Duration;
@@ -67,7 +67,10 @@ async fn s3_put_get_list_succeed_during_decommission_and_rebalance() -> TestResu
assert_inventory(&live, &bucket, &inventory).await?;
let rebalance_id = start_rebalance(&dist.cluster).await?;
wait_for_rebalance_running_with_progress(&dist.cluster, &rebalance_id, Duration::from_secs(30)).await?;
// The status API reads persisted progress, whose first periodic save is
// after 30 seconds. A shorter run can remain at zero until completion.
// Require Started around the S3 operations and nonzero progress at completion.
wait_for_rebalance_active(&dist.cluster, &rebalance_id, Duration::from_secs(30)).await?;
retrying_put(
&live,
&bucket,
@@ -84,11 +87,26 @@ async fn s3_put_get_list_succeed_during_decommission_and_rebalance() -> TestResu
Duration::from_secs(30),
)
.await?;
let listed = live.list_objects_v2().bucket(&bucket).send().await?;
assert!(
listed
.contents()
.iter()
.any(|object| object.key() == Some("during-rebalance.bin")),
"list during rebalance missed the newly written key"
);
let status = rebalance_status_json(&dist.cluster).await?;
if !rebalance_running_with_progress(&status, &rebalance_id)? {
if !rebalance_active(&status, &rebalance_id)? {
return Err(format!("rebalance did not remain active across the S3 operations: {status}").into());
}
wait_for_rebalance_complete(&dist.cluster, &rebalance_id, Duration::from_secs(180)).await?;
assert_inventory(&dist.client(1)?, &bucket, &inventory).await?;
let after = dist.client(1)?;
assert_inventory(&after, &bucket, &inventory).await?;
for (key, body) in [
("during-decommission.bin", b"written-while-decommissioning".as_slice()),
("during-rebalance.bin", b"written-while-rebalancing".as_slice()),
] {
retrying_get_equals(&after, &bucket, key, body, Duration::from_secs(30)).await?;
}
Ok(())
}
@@ -1569,7 +1569,9 @@ mod tests {
} else {
if matches!(
scenario,
InterruptionScenario::BackgroundTargetRestart | InterruptionScenario::BackgroundTargetRestartEc84
InterruptionScenario::BackgroundTargetRestart
| InterruptionScenario::BackgroundTargetRestartEc84
| InterruptionScenario::BackgroundCoordinatorRestart
) {
cluster.stop_node_gracefully(interruption_node).await?;
} else {
@@ -1767,33 +1769,22 @@ mod tests {
let task_status_body = signed_admin_post(&task_status_url, None, &cluster.access_key, &cluster.secret_key).await?;
let task_status: serde_json::Value = serde_json::from_str(&task_status_body)
.map_err(|err| format!("heal task status is not JSON ({err}): {task_status_body}"))?;
if task_status["summary"].as_str() != Some("finished") {
return Err(format!("heal data rebuilt but task did not finish successfully: {task_status}").into());
}
if interruption_node == 0 {
// Admin tasks are process-local. Physical and queue convergence
// above establish recovery; a lost task must not report success.
assert_eq!(
task_status["summary"].as_str(),
Some("notFound"),
"interrupted task status: {task_status}"
);
assert_eq!(
task_status["detail"].as_str(),
Some("heal task not found or expired"),
"interrupted admin task must be explicitly unavailable: {task_status}"
);
// Restart recovery must finish the original durable root request.
info!(
event = "heal_interruption_recovered",
component = "e2e_test",
subsystem = "heal",
interruption_node,
interruption_kind,
task_state = "not_found",
"Physical recovery completed after coordinator restart"
task_state = "finished",
"Original root heal completed after coordinator restart"
);
return Ok(());
}
if task_status["summary"].as_str() != Some("finished") {
return Err(format!("heal data rebuilt but task did not finish successfully: {task_status}").into());
}
if let Some(evidence_context) = evidence_run {
let restarted_pid = cluster.nodes[1].process.as_ref().ok_or("restarted target is absent")?.id();
@@ -109,10 +109,10 @@ async fn test_kms_key_directory_unavailable() -> Result<(), Box<dyn std::error::
.await;
let unavailable_error = put_result2.expect_err("a missing Local KMS key directory must reject encrypted writes");
assert_eq!(unavailable_error.raw_response().map(|response| response.status().as_u16()), Some(500));
assert_eq!(unavailable_error.raw_response().map(|response| response.status().as_u16()), Some(503));
assert_eq!(
unavailable_error.as_service_error().and_then(ProvideErrorMetadata::code),
Some("InternalError")
Some("ServiceUnavailable")
);
let unavailable_absence = s3_client
.get_object()
+435 -13
View File
@@ -2606,6 +2606,24 @@ where
usize::try_from(size).unwrap_or_default()
}
fn is_decommission_set_local_usage_cache(bucket: &str, object: &str) -> bool {
if bucket != RUSTFS_META_BUCKET {
return false;
}
let Some(path) = object
.strip_prefix(BUCKET_META_PREFIX)
.and_then(|path| path.strip_prefix('/'))
else {
return false;
};
let name = match path.rsplit_once('/') {
Some((bucket, name)) if !bucket.is_empty() && !bucket.contains('/') && bucket != "." && bucket != ".." => name,
Some(_) => return false,
None => path,
};
name.strip_suffix(".bkp").unwrap_or(name) == DATA_USAGE_CACHE_NAME
}
fn with_decommission_entry_context<E: Display>(stage: &str, bucket: &str, object: &str, err: E) -> Error {
Error::other(format!("decommission entry {stage} failed for bucket {bucket} object {object}: {err}"))
}
@@ -13285,6 +13303,12 @@ impl ECStore {
);
return Ok(DecommissionEntryAttemptOutcome::Complete);
}
// Scanner caches describe their own erasure set and are rebuilt there.
// Copying one onto another set can overwrite unrelated cache contents or
// leave an unresolvable target-capacity intent after a conditional PUT.
if is_decommission_set_local_usage_cache(&bucket, &entry.name) {
return Ok(DecommissionEntryAttemptOutcome::Complete);
}
let durable_ilm_record = if bucket == RUSTFS_META_BUCKET {
classify_durable_ilm_record(&entry.name)
.map_err(|err| with_decommission_entry_context("durable_ilm_namespace", &bucket, &entry.name, err))?
@@ -13842,7 +13866,7 @@ impl ECStore {
let bucket = bucket.clone();
let rd = match set
let read_result = set
.get_object_reader(
bucket.as_str(),
&encode_dir_object(&version.name),
@@ -13850,8 +13874,11 @@ impl ECStore {
HeaderMap::new(),
&decommission_object_migration_read_opts(version_id.clone()),
)
.await
{
.await;
#[cfg(test)]
let read_result =
decommission_test_wrap_result("object_read", &bucket, &version.name, version_attempt, read_result);
let rd = match read_result {
Ok(rd) => rd,
Err(err) => {
if is_err_object_not_found(&err) || is_err_version_not_found(&err) {
@@ -13860,15 +13887,6 @@ impl ECStore {
break;
}
if !ignore {
//
if bucket == RUSTFS_META_BUCKET && version.name.contains(DATA_USAGE_CACHE_NAME) {
ignore = true;
error!("decommission_pool: ignore data usage cache {}", &version.name);
break;
}
}
failure = true;
if version_attempt == DECOMMISSION_VERSION_COPY_ATTEMPTS {
error!(
@@ -16836,7 +16854,7 @@ impl ECStore {
return;
}
if bucket_name == RUSTFS_META_BUCKET && entry.name.contains(DATA_USAGE_CACHE_NAME) {
if is_decommission_set_local_usage_cache(&bucket_name, &entry.name) {
return;
}
@@ -17154,6 +17172,361 @@ mod tests {
use crate::storage_api_contracts::multipart::MultipartOperations as _;
use serde::Serialize;
#[test]
fn decommission_set_local_usage_cache_classification_is_exact() {
for object in [
"buckets/.usage-cache.bin",
"buckets/.usage-cache.bin.bkp",
"buckets/photos/.usage-cache.bin",
"buckets/photos/.usage-cache.bin.bkp",
] {
assert!(is_decommission_set_local_usage_cache(RUSTFS_META_BUCKET, object), "{object}");
assert!(!is_decommission_set_local_usage_cache("user-bucket", object), "{object}");
}
for object in [
"buckets/.usage.v2.json",
"buckets/.usage.v2.json.bkp",
"buckets/.usage-cache.bin.extra",
"buckets/.usage-cache.bin.bkp.extra",
"buckets/prefix.usage-cache.bin",
"buckets/photos/.usage-cache.bin.bkp.bkp",
"buckets/photos/nested/.usage-cache.bin",
"buckets//.usage-cache.bin",
"buckets/../.usage-cache.bin",
"buckets/./.usage-cache.bin",
"buckets/.usage-cache.bin/child",
"config/.usage-cache.bin",
"buckets-other/.usage-cache.bin",
".usage-cache.bin",
] {
assert!(!is_decommission_set_local_usage_cache(RUSTFS_META_BUCKET, object), "{object}");
}
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_keeps_set_local_usage_caches_out_of_target_capacity() {
use crate::object_api::PutObjReader;
use tokio::io::AsyncReadExt as _;
// Keep the scenario's large setup and migration futures off the test
// future so ordinary metadata I/O retains the default thread stack.
let (_temp_dirs, store, _other_store) =
Box::pin(crate::services::rebalance::test_two_pool_stores_with_isolated_node_contexts(None)).await;
let user_bucket = "decommission-usage-cache-control";
Box::pin(store.make_bucket(user_bucket, &MakeBucketOptions::default()))
.await
.expect("create the ordinary-object control bucket");
let incarnation = Box::pin(store.bucket_incarnation_id(user_bucket))
.await
.expect("control bucket incarnation");
let source_time = OffsetDateTime::now_utc();
let cache_objects = [
"buckets/.usage-cache.bin",
"buckets/.usage-cache.bin.bkp",
"buckets/photos/.usage-cache.bin",
"buckets/photos/.usage-cache.bin.bkp",
];
let conflict_object = "buckets/.usage-cache.bin.conflict";
let source_read_failure_object = "buckets/.usage-cache.bin.read-error";
let source_body = b"source set cache";
let target_body = b"independent older target set cache";
for object in cache_objects.into_iter().chain([conflict_object]) {
for (pool_index, body, mod_time) in [
(0, source_body.as_slice(), source_time),
(1, target_body.as_slice(), source_time - Duration::seconds(1)),
] {
// Scanner cache persistence writes directly to its own set.
store.pools[pool_index]
.get_disks_by_key(object)
.put_object(
RUSTFS_META_BUCKET,
object,
&mut PutObjReader::from_vec(body.to_vec()),
&ObjectOptions {
mod_time: Some(mod_time),
..Default::default()
},
)
.await
.expect("seed distinct native set-local objects");
}
}
let controls = [
(RUSTFS_META_BUCKET, "buckets/.usage.v2.json"),
(RUSTFS_META_BUCKET, "buckets/photos/.usage-cache.bin.extra"),
(user_bucket, "ordinary-object"),
(user_bucket, "buckets/.usage-cache.bin"),
];
for (bucket, object) in controls {
store.pools[0]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(b"ordinary object contents".to_vec()),
&ObjectOptions {
expected_bucket_incarnation_id: (bucket == user_bucket).then_some(incarnation),
..Default::default()
},
)
.await
.expect("seed a control that must migrate");
}
store.pools[0]
.put_object(
RUSTFS_META_BUCKET,
source_read_failure_object,
&mut PutObjReader::from_vec(source_body.to_vec()),
&ObjectOptions::default(),
)
.await
.expect("seed a similarly named object whose source read will fail");
let layout = DecommissionErasureLayout { data: 1, parity: 0 };
set_decommission_capacity_info_overrides_for_test(
store.id,
vec![vec![
DecommissionPoolCapacityInfo::for_test(0, layout, 0, 16_384, 16_384),
DecommissionPoolCapacityInfo::for_test(1, layout, 131_072, 131_072, 0),
]],
);
Box::pin(store.save_current_pool_meta_for_decommission_start(&[0], Vec::new()))
.await
.expect("activate the decommission capacity reservation");
for object in cache_objects {
Box::pin(store.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
RUSTFS_META_BUCKET.to_string(),
store.pools[0].get_disks_by_key(object),
))
.await
.expect("set-local cache must not enter cross-pool migration");
for (pool_index, expected) in [(0, source_body.as_slice()), (1, target_body.as_slice())] {
let mut reader = store.pools[pool_index]
.get_disks_by_key(object)
.get_object_reader(RUSTFS_META_BUCKET, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("each set must retain its own cache");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read retained cache");
assert_eq!(actual, expected, "pool {pool_index}, {object}");
}
let meta = store.pool_meta.read().await;
let info = meta.pools[0].decommission.as_ref().expect("decommission progress");
assert_eq!((info.items_decommissioned, info.items_decommission_failed), (0, 0));
assert_eq!((info.bytes_done, info.bytes_failed), (0, 0));
let reservation = info.capacity_reservation.as_ref().expect("capacity reservation");
assert_eq!(reservation.pending_target_physical_bytes, 0);
assert_eq!(reservation.consumed_target_physical_bytes, 0);
assert!(reservation.targets.iter().all(|target| target.pending_mutation_id.is_none()));
}
let mut persisted = PoolMeta::default();
Box::pin(persisted.load_no_lock_from_replicas(store.pools.clone()))
.await
.expect("reload durable capacity intents after cache entries");
let reservation = persisted.pools[0]
.decommission
.as_ref()
.and_then(|info| info.capacity_reservation.as_ref())
.expect("durable reservation");
assert_eq!(reservation.pending_target_physical_bytes, 0);
assert!(reservation.targets.iter().all(|target| target.pending_mutation_id.is_none()));
let injected_reads = Arc::new(AtomicUsize::new(0));
let observed_reads = Arc::clone(&injected_reads);
let read_fault = DecommissionTestFaultGuard::install(Arc::new(move |stage, bucket, object, _, success| {
if stage == "object_read" && bucket == RUSTFS_META_BUCKET && object == source_read_failure_object && success {
observed_reads.fetch_add(1, Ordering::SeqCst);
return true;
}
false
}));
Box::pin(store.decommission_entry_for_test(
0,
MetaCacheEntry {
name: source_read_failure_object.to_string(),
..Default::default()
},
RUSTFS_META_BUCKET.to_string(),
store.pools[0].get_disks_by_key(source_read_failure_object),
))
.await
.expect("entry must record the non-NotFound source read failure");
drop(read_fault);
assert_eq!(injected_reads.load(Ordering::SeqCst), DECOMMISSION_VERSION_COPY_ATTEMPTS);
{
let meta = store.pool_meta.read().await;
let info = meta.pools[0].decommission.as_ref().expect("source read failure progress");
assert_eq!((info.items_decommissioned, info.items_decommission_failed), (0, 1));
assert_eq!(info.bytes_failed, source_body.len());
assert_eq!(
info.capacity_reservation
.as_ref()
.expect("reservation")
.pending_target_physical_bytes,
0
);
}
let mut retained = store.pools[0]
.get_object_reader(
RUSTFS_META_BUCKET,
source_read_failure_object,
None,
HeaderMap::new(),
&ObjectOptions::default(),
)
.await
.expect("source read failure must retain the source");
let mut retained_body = Vec::new();
retained
.stream
.read_to_end(&mut retained_body)
.await
.expect("read retained source");
assert_eq!(retained_body, source_body);
drop(retained);
let target_err = store.pools[1]
.get_object_info(RUSTFS_META_BUCKET, source_read_failure_object, &ObjectOptions::default())
.await
.expect_err("failed source read must not create a target object");
assert!(is_err_object_not_found(&target_err), "unexpected target state: {target_err:?}");
for (bucket, object) in controls {
Box::pin(store.decommission_entry_for_test(
0,
MetaCacheEntry {
name: object.to_string(),
..Default::default()
},
bucket.to_string(),
store.pools[0].get_disks_by_key(object),
))
.await
.expect("ordinary and similarly named objects must migrate");
let mut reader = store.pools[1]
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("control must exist on the target");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read migrated control");
assert_eq!(actual, b"ordinary object contents", "{bucket}/{object}");
let err = store.pools[0]
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect_err("migrated control must be removed from the source");
assert!(is_err_object_not_found(&err), "{bucket}/{object}: {err:?}");
}
Box::pin(store.decommission_entry_for_test(
0,
MetaCacheEntry {
name: conflict_object.to_string(),
..Default::default()
},
RUSTFS_META_BUCKET.to_string(),
store.pools[0].get_disks_by_key(conflict_object),
))
.await
.expect("entry must record a real conditional-copy failure");
let meta = store.pool_meta.read().await;
let info = meta.pools[0].decommission.as_ref().expect("final progress");
assert_eq!(info.items_decommissioned, controls.len());
assert_eq!(
info.items_decommission_failed, 2,
"similar names must not hide read or migration failures"
);
assert_eq!(info.bytes_failed, source_body.len() * 2);
drop(meta);
for (pool_index, expected) in [(0, source_body.as_slice()), (1, target_body.as_slice())] {
let mut reader = store.pools[pool_index]
.get_object_reader(RUSTFS_META_BUCKET, conflict_object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("failed migration must preserve both objects");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read conflict object");
assert_eq!(actual, expected);
}
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_final_sweep_excludes_only_set_local_usage_caches() {
use crate::object_api::PutObjReader;
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_two_pool_stores_with_isolated_node_contexts(None).await;
for object in [
"buckets/.usage-cache.bin",
"buckets/.usage-cache.bin.bkp",
"buckets/photos/.usage-cache.bin",
"buckets/photos/.usage-cache.bin.bkp",
] {
store.pools[0]
.get_disks_by_key(object)
.put_object(
RUSTFS_META_BUCKET,
object,
&mut PutObjReader::from_vec(b"set-local cache".to_vec()),
&ObjectOptions::default(),
)
.await
.expect("seed each supported set-local cache path");
}
let layout = DecommissionErasureLayout { data: 1, parity: 0 };
set_decommission_capacity_info_overrides_for_test(
store.id,
vec![vec![
DecommissionPoolCapacityInfo::for_test(0, layout, 0, 16_384, 16_384),
DecommissionPoolCapacityInfo::for_test(1, layout, 131_072, 131_072, 0),
]],
);
store
.save_current_pool_meta_for_decommission_start(&[0], Vec::new())
.await
.expect("activate the final-sweep generation");
let generation = store.active_decommission_generation(0).await.expect("active generation");
store
.check_after_decommission(0, &CancellationToken::new(), generation)
.await
.expect("the four set-local cache forms must not block the final sweep");
for object in [
"buckets/.usage-cache.bin.extra",
"buckets/photos/.usage-cache.bin.bkp.extra",
"buckets/.usage.v2.json",
] {
let source_set = store.pools[0].get_disks_by_key(object);
source_set
.put_object(
RUSTFS_META_BUCKET,
object,
&mut PutObjReader::from_vec(b"unmigrated ordinary metadata".to_vec()),
&ObjectOptions::default(),
)
.await
.expect("seed ordinary metadata that must prevent completion");
let err = store
.check_after_decommission(0, &CancellationToken::new(), generation)
.await
.expect_err("a remaining similar name or global usage snapshot must block completion");
assert!(err.to_string().contains("after decommissioning"), "unexpected final-sweep error: {err:?}");
assert!(err.to_string().contains(object), "the final sweep must identify {object}: {err:?}");
source_set
.delete_object(RUSTFS_META_BUCKET, object, ObjectOptions::default())
.await
.expect("remove only the ordinary-metadata control before the next sweep");
}
store
.check_after_decommission(0, &CancellationToken::new(), generation)
.await
.expect("only the four set-local caches remain after removing the controls");
}
#[test]
fn pool_activation_fleet_proof_error_classifier_matches_only_retryable_proof_failures() {
assert!(is_pool_activation_fleet_proof_error(&Error::other(POOL_ACTIVATION_FLEET_PROOF_REQUIRED)));
@@ -19865,6 +20238,55 @@ mod tests {
assert!(!is_decommission_copy_cleanup_safe_error(&wrap(Error::SlowDown)));
}
#[test]
fn decommission_target_gate_retry_recognizes_multipart_part_errors() {
let wrap = |inner: Error| {
data_movement::data_movement_part_stage_error_for_test(
"decommission_object",
"put_object_part",
"bucket-a",
"object-a",
1,
inner,
)
};
let gate_busy_message =
format!("{DECOMMISSION_CAPACITY_TARGET_GATE_BUSY_PREFIX}7{DECOMMISSION_CAPACITY_TARGET_GATE_BUSY_SUFFIX}");
let wrapped = wrap(decommission_capacity_blocked_error(&gate_busy_message));
assert!(is_decommission_capacity_target_gate_busy(&wrapped));
assert_eq!(decommission_capacity_target_gate_busy_index(&wrapped), Some(7));
assert_eq!(
wrapped.to_string(),
format!(
"Io error: decommission_object: put_object_part failed for bucket-a/object-a part 1: {}",
decommission_capacity_blocked_error(&gate_busy_message)
)
);
for unrelated in [
Error::SlowDown,
Error::DiskFull,
decommission_capacity_blocked_error("target capacity is exhausted"),
Error::other(gate_busy_message),
] {
let wrapped = wrap(unrelated);
assert!(!is_decommission_capacity_target_gate_busy(&wrapped));
assert_eq!(decommission_capacity_target_gate_busy_index(&wrapped), None);
}
for missing_target in [
Error::FileNotFound,
Error::ObjectNotFound("bucket-a".to_string(), "object-a".to_string()),
Error::VersionNotFound("bucket-a".to_string(), "object-a".to_string(), "version-a".to_string()),
] {
assert!(is_decommission_copy_cleanup_safe_error(&missing_target));
assert!(
!is_decommission_copy_cleanup_safe_error(&wrap(missing_target)),
"a missing target part must never authorize source cleanup"
);
}
assert!(is_decommission_target_capacity_error(&wrap(Error::DiskFull)));
}
#[test]
fn decommission_target_capacity_error_accepts_wrapped_capacity_errors() {
let disk_full = Error::other(format!("decommission_object: put_object failed for bucket/object: {}", Error::DiskFull));
+29 -4
View File
@@ -1521,9 +1521,27 @@ fn data_movement_part_stage_error(
bucket: &str,
object: &str,
part_number: usize,
err: impl std::fmt::Display,
err: Error,
) -> Error {
Error::other(format!("{op_label}: {stage} failed for {bucket}/{object} part {part_number}: {err}"))
let rendered = format!("{op_label}: {stage} failed for {bucket}/{object} part {part_number}: {err}");
if matches!(&err, Error::DecommissionCapacityBlocked { .. }) {
return data_movement_context_error(rendered, err);
}
// A missing target part is not evidence that the source can be deleted.
// Keep other part errors opaque to the source-cleanup classifiers.
Error::other(rendered)
}
#[cfg(test)]
pub(crate) fn data_movement_part_stage_error_for_test(
op_label: &str,
stage: &str,
bucket: &str,
object: &str,
part_number: usize,
err: Error,
) -> Error {
data_movement_part_stage_error(op_label, stage, bucket, object, part_number, err)
}
fn is_data_movement_part_read_error(err: &Error) -> bool {
@@ -2428,8 +2446,15 @@ mod tests {
let err =
data_movement_part_stage_error("rebalance_object", "put_object_part", "bucket-a", "object-a", 7, Error::SlowDown);
let message = err.to_string();
assert!(message.contains("rebalance_object: put_object_part failed for bucket-a/object-a part 7"));
assert!(message.contains(Error::SlowDown.to_string().as_str()));
assert_eq!(
message,
Error::other(format!(
"rebalance_object: put_object_part failed for bucket-a/object-a part 7: {}",
Error::SlowDown
))
.to_string()
);
assert!(data_movement_stage_source(&err).is_none());
}
#[test]
+248 -9
View File
@@ -164,6 +164,16 @@ pub fn max_keys_plus_one(max_keys: i32, add_one: bool) -> i32 {
max_keys
}
fn list_versions_scan_limit(max_keys: i32, has_version_marker: bool) -> i32 {
if max_keys <= 0 {
return 0;
}
// The marker object's versions may all be filtered out after gathering.
// Reserve its raw entry in addition to the next-page lookahead entry.
max_keys_plus_one(max_keys, true) + i32::from(has_version_marker)
}
#[derive(Debug, Clone, Copy, Eq, PartialEq)]
enum GatherResultsState {
LimitReached,
@@ -2139,15 +2149,19 @@ fn build_list_versions_next_marker(
// here; advertise it as the literal `null` marker so a resumed listing
// parses it back to `VersionMarker::Null` instead of a nil UUID that
// `find_version_index` can never match (issue #6745).
(
Some(append_list_cache_id_to_marker(last.name.clone(), cache_id)),
let version_marker = if last.is_dir && last.mod_time.is_none() {
// A CommonPrefix has no version to resume; a version marker would
// make the next page include this same prefix again.
None
} else {
Some(
last.version_id
.filter(|v| !v.is_nil())
.map(|v| v.to_string())
.unwrap_or_else(|| "null".to_string()),
),
)
)
};
(Some(append_list_cache_id_to_marker(last.name.clone(), cache_id)), version_marker)
} else if let Some(last_prefix) = prefixes.last() {
(Some(append_list_cache_id_to_marker(last_prefix.clone(), cache_id)), None)
} else {
@@ -2866,6 +2880,20 @@ fn listing_entries_supplement_target(
return None;
}
if let Some(directory) = entries.0.iter().flatten().find(|entry| entry.is_dir()) {
let directory_copies = entries
.0
.iter()
.flatten()
.filter(|entry| entry.is_dir() && entry.name == directory.name)
.count();
// A committed child may have some of its directory copies only on
// fallback disks, just like object metadata in a partial primary sample.
if directory_copies < resolver.dir_quorum {
return Some(directory.name.clone());
}
}
for (idx, entry) in entries.0.iter().enumerate() {
let Some(entry) = entry.as_ref().filter(|entry| entry.is_object()) else {
continue;
@@ -4018,8 +4046,7 @@ impl ECStore {
None
};
let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) };
// Always request max_keys + 1 to detect if there are more results
let effective_max_keys = list_versions_scan_limit(max_keys, has_version_marker);
let mut opts = ListPathOptions {
bucket: bucket.to_owned(),
prefix: prefix.to_owned(),
@@ -5325,7 +5352,7 @@ impl Sets {
None
};
let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) };
let effective_max_keys = list_versions_scan_limit(max_keys, has_version_marker);
let mut opts = ListPathOptions {
bucket: bucket.to_owned(),
prefix: prefix.to_owned(),
@@ -6034,7 +6061,7 @@ impl SetDisks {
let has_version_marker = version_marker.is_some();
let version_marker = version_marker.map(parse_version_marker).transpose()?;
let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) };
let effective_max_keys = list_versions_scan_limit(max_keys, has_version_marker);
let mut opts = ListPathOptions {
bucket: bucket.to_owned(),
prefix: prefix.to_owned(),
@@ -6248,7 +6275,7 @@ impl SetDisks {
None
};
let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) };
let effective_max_keys = list_versions_scan_limit(max_keys, has_version_marker);
let mut opts = ListPathOptions {
bucket: bucket.to_owned(),
prefix: prefix.to_owned(),
@@ -7441,6 +7468,153 @@ mod test {
assert!(cancel.is_cancelled());
}
#[test]
fn list_versions_pagination_scan_limit_boundaries() {
for has_version_marker in [false, true] {
assert_eq!(super::list_versions_scan_limit(-1, has_version_marker), 0);
assert_eq!(super::list_versions_scan_limit(0, has_version_marker), 0);
let marker_slot = i32::from(has_version_marker);
assert_eq!(super::list_versions_scan_limit(1, has_version_marker), 2 + marker_slot);
assert_eq!(super::list_versions_scan_limit(MAX_OBJECT_LIST, has_version_marker), 1001 + marker_slot);
assert_eq!(super::list_versions_scan_limit(i32::MAX, has_version_marker), 1001 + marker_slot);
}
}
#[tokio::test]
async fn list_versions_pagination_does_not_require_an_empty_final_page() {
use crate::bucket::metadata_sys::{init_bucket_metadata_sys, test_support::isolated_store_over_temp_disks};
use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions};
let (dirs, store) = isolated_store_over_temp_disks().await;
let bucket = "version-pagination-bucket";
init_bucket_metadata_sys(store.clone(), Vec::new()).await;
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("pagination bucket should be created");
let mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
for kind in ["objects", "deletes", "null", "mixed", "delimiter"] {
let count = if kind == "mixed" { 5 } else { 10 };
let mut expected = Vec::new();
for index in 0..count {
let name = if kind == "delimiter" && index % 2 == 1 {
format!("{kind}/testobject-{index:02}/child")
} else {
format!("{kind}/testobject-{index:02}")
};
let entry = match kind {
"deletes" => test_delete_marker_meta_entry(&name, mod_time),
"null" => test_object_meta_entry(&name),
"mixed" => test_object_with_delete_marker_meta_entry(&name, mod_time, mod_time + time::Duration::SECOND),
_ => test_object_meta_entry_with_erasure_versions(&name, &[(mod_time, "etag", 2, 2)]),
};
for dir in &dirs {
let object_dir = dir.path().join(bucket).join(&name);
tokio::fs::create_dir_all(&object_dir)
.await
.expect("pagination object directory should be created");
tokio::fs::write(object_dir.join(STORAGE_FORMAT_FILE), &entry.metadata)
.await
.expect("pagination metadata should be written");
}
if kind == "delimiter" && index % 2 == 1 {
expected.push((name.trim_end_matches("child").to_owned(), None, false));
} else {
let versions = entry.file_info_versions(bucket).expect("fixture versions should decode");
expected.extend(
versions
.versions
.iter()
.map(|version| (name.clone(), version.version_id, version.deleted)),
);
}
}
let prefix = format!("{kind}/");
let delimiter = (kind == "delimiter").then(|| "/".to_owned());
// Exercise each public/internal entry point with the reported page size.
// The store entry point also covers exact and one-over limit boundaries.
for (layer, max_keys) in [(0, 0), (0, 1), (0, 5), (0, 9), (0, 10), (0, 11), (1, 5), (2, 5), (3, 5)] {
if layer == 3 && delimiter.is_some() {
continue;
}
let mut marker = None;
let mut version_marker = None;
let expected_pages = if max_keys == 0 {
1
} else {
10usize.div_ceil(usize::try_from(max_keys).expect("positive page size"))
};
let mut actual = Vec::new();
for page in 0..expected_pages {
let result = match layer {
0 => {
store
.clone()
.inner_list_object_versions(bucket, &prefix, marker, version_marker, delimiter.clone(), max_keys)
.await
}
1 => {
store.pools[0]
.clone()
.inner_list_object_versions(bucket, &prefix, marker, version_marker, delimiter.clone(), max_keys)
.await
}
2 => {
store.pools[0].disk_set[0]
.clone()
.inner_list_object_versions(bucket, &prefix, marker, version_marker, delimiter.clone(), max_keys)
.await
}
_ => {
store.pools[0].disk_set[0]
.clone()
.inner_list_object_versions_for_recursive_delete(
bucket,
&prefix,
marker,
version_marker,
max_keys,
)
.await
}
}
.expect("version page should list successfully");
let page_size = usize::try_from(max_keys).expect("nonnegative page size");
assert_eq!(result.objects.len() + result.prefixes.len(), (10 - page * page_size).min(page_size));
let has_more = page + 1 < expected_pages;
assert_eq!(result.is_truncated, has_more, "{kind}, layer {layer}, max_keys {max_keys}, page {page}");
assert_eq!(
result.next_marker.is_some(),
has_more,
"key marker must exist only when another page exists"
);
if !has_more {
assert!(
result.next_version_idmarker.is_none(),
"the final page must not advertise a version marker"
);
}
actual.extend(
result
.objects
.into_iter()
.map(|object| (object.name, object.version_id, object.delete_marker)),
);
actual.extend(result.prefixes.into_iter().map(|prefix| (prefix, None, false)));
marker = result.next_marker;
version_marker = result.next_version_idmarker;
}
// Objects and CommonPrefixes are serialized separately; compare their
// identities without relying on their relative position in the response.
actual.sort();
let mut expected = if max_keys == 0 { Vec::new() } else { expected.clone() };
expected.sort();
assert_eq!(actual, expected, "{kind}, layer {layer}, max_keys {max_keys}");
}
}
}
#[test]
fn version_marker_is_applied_only_when_key_marker_entry_is_present() {
let version_marker = Some(VersionMarker::Null);
@@ -9448,6 +9622,71 @@ mod test {
assert!(supplemented.is_latest_delete_marker());
}
#[tokio::test]
async fn latest_listing_supplement_checks_fallback_disks_for_common_prefix_quorum() {
let mut fallback_disks = Vec::new();
let mut fallback_tempdirs = Vec::new();
for index in 0..4 {
let tempdir = tempfile::tempdir().expect("fallback tempdir should be created");
let endpoint = Endpoint::try_from(tempdir.path().to_str().expect("fallback path should be utf8"))
.expect("fallback endpoint should parse");
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("fallback disk should be created");
disk.make_volume("bucket").await.expect("fallback bucket should be created");
for copies in [3, 4] {
if index < copies {
let object = format!("quux-{copies}/thud");
let entry = test_object_meta_entry(&object);
disk.write_all("bucket", &format!("{object}/{STORAGE_FORMAT_FILE}"), bytes::Bytes::from(entry.metadata))
.await
.expect("fallback child metadata should be written");
}
}
fallback_disks.push(disk);
fallback_tempdirs.push(tempdir);
}
let supplement = ListingSupplement::new(
ListingSupplementOptions {
bucket: "bucket".to_owned(),
path: String::new(),
recursive: false,
incl_deleted: false,
skip_hidden_prefix_check: false,
filter_prefix: None,
forward_to: None,
per_disk_limit: 100,
skip_total_timeout: true,
walkdir_timeout: None,
walkdir_stall_timeout: None,
},
Arc::new(fallback_disks),
FallbackClaimTracker::default(),
);
// A 16-drive EC:4 set asks 12 primary disks. A committed write may
// exist on eight primary disks and all four remaining fallback disks.
let resolver = list_metadata_resolution_params("bucket".to_owned(), 4, 12, false, 0);
for fallback_copies in [3, 4] {
let prefix = format!("quux-{fallback_copies}/");
let mut primary = vec![Some(test_dir_meta_entry(&prefix)); 8];
primary.extend([None, None, None, None]);
let entry =
resolve_listing_entries_with_supplement(MetaCacheEntries(primary), resolver.clone(), true, supplement.clone())
.await;
assert_eq!(
entry.map(|entry| entry.name),
(fallback_copies == 4).then_some(prefix),
"the common prefix needs all twelve copies, including fallback disks"
);
}
}
#[test]
fn latest_listing_supplement_keeps_a_subquorum_delete_marker_hidden() {
let object_mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
+130 -17
View File
@@ -841,6 +841,10 @@ pub struct HealManager {
replacement_recovery_anchors: Arc<std::sync::Mutex<HashMap<String, String>>>,
/// Set IDs whose durable replacement metadata is corrupt or conflicting.
replacement_recovery_blocked_sets: Arc<std::sync::Mutex<HashSet<String>>>,
/// Durable handoff of interrupted administrator root traversals.
root_recovery: Arc<root_recovery::RootHealRecovery>,
/// Keep forceStart's cancellation side effects inside the shutdown fence.
force_start_shutdown: Mutex<()>,
/// Storage layer interface
storage: Arc<dyn HealStorageAPI>,
/// Cancel token
@@ -876,6 +880,7 @@ struct HealQueueContext<'a> {
retrying_heals: &'a Arc<Mutex<HashMap<String, RetryingHeal>>>,
mrf_repair_notice_targets: &'a Arc<StdMutex<HashMap<String, Vec<MrfRepairNoticeTarget>>>>,
replacement_recovery_anchors: &'a Arc<std::sync::Mutex<HashMap<String, String>>>,
root_recovery: &'a Arc<root_recovery::RootHealRecovery>,
config: &'a Arc<RwLock<HealConfig>>,
statistics: &'a Arc<RwLock<HealStatistics>>,
storage: &'a Arc<dyn HealStorageAPI>,
@@ -1377,6 +1382,8 @@ impl HealManager {
mrf_repair_notice_targets: Arc::new(StdMutex::new(HashMap::new())),
replacement_recovery_anchors: Arc::new(std::sync::Mutex::new(HashMap::new())),
replacement_recovery_blocked_sets: Arc::new(std::sync::Mutex::new(HashSet::new())),
root_recovery: Arc::new(root_recovery::RootHealRecovery::default()),
force_start_shutdown: Mutex::new(()),
storage,
cancel_token: CancellationToken::new(),
statistics: Arc::new(RwLock::new(HealStatistics::new())),
@@ -1412,6 +1419,23 @@ impl HealManager {
"Heal manager starting"
);
// Restore graceful-shutdown root responsibilities before automatic
// repair can admit overlapping work.
if let Err(error) = self.replay_root_heals().await {
// A missing owner or invalid root record must not block existing
// replacement recovery. Keep its file for a later restart after
// the owner is readable or the record has been repaired.
warn!(
target: "rustfs::heal::manager",
event = EVENT_HEAL_MANAGER_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_MANAGER,
state = "root_recovery_deferred",
error = %error,
"Root heal restart recovery deferred"
);
}
// start scheduler
self.start_scheduler().await?;
@@ -1449,6 +1473,7 @@ impl HealManager {
/// Stop HealManager
pub async fn stop(&self) -> Result<()> {
let _force_start_guard = self.force_start_shutdown.lock().await;
info!(
target: "rustfs::heal::manager",
event = EVENT_HEAL_MANAGER_STATE,
@@ -1458,11 +1483,39 @@ impl HealManager {
"Heal manager stopping"
);
// cancel all tasks
self.cancel_token.cancel();
// wait for all tasks to complete
// Keep scheduler, cancellation, and retry ownership stable until every
// unfinished root traversal has a durable successor. A failed write
// must leave the manager running and the shutdown marker unclean.
let mut active_heals = self.active_heals.lock().await;
let queue = self.heal_queue.lock().await;
let retrying = self.retrying_heals.lock().await;
for task in active_heals.values() {
if root_recovery::is_root_heal(&task.heal_type, task.source) {
if task.get_status().await == HealTaskStatus::Completed {
self.root_recovery.remove(&task.id, &task.heal_type, task.source).await?;
} else {
let mut request = match task.retry_request_with_remaining_timeout().await {
Ok(request) => request,
Err(Error::TaskTimeout) => {
let mut request = task.retry_request();
request.options.timeout = Some(Duration::ZERO);
request
}
Err(error) => return Err(error),
};
request.retry_attempts = task.retry_attempts;
self.root_recovery.persist(&request).await?;
}
}
}
for request in queue.requests().chain(retrying.values().map(|retrying| &retrying.request)) {
self.root_recovery.persist(request).await?;
}
self.cancel_token.cancel();
drop(retrying);
drop(queue);
// cancel active workers after the durable handoff
for task in active_heals.values() {
if let Err(e) = task.cancel().await {
warn!(
@@ -1589,6 +1642,17 @@ impl HealManager {
let admission_start = Instant::now();
let source = request.source;
let force_start = request.force_start;
// A forceStart must not retire an old durable owner if shutdown will
// reject its replacement. Hold the same gate through final admission.
let _force_start_guard = if source == HealRequestSource::Admin && force_start {
let guard = self.force_start_shutdown.lock().await;
if self.cancel_token.is_cancelled() {
return Err(Error::Other("Heal manager is stopping".to_string()));
}
Some(guard)
} else {
None
};
// HS-06 forceStart semantics (admin only): MinIO stops the old task
// first and then starts the new one. Cancel any active admin task
// overlapping this request's path before entering admission, so the
@@ -1596,7 +1660,9 @@ impl HealManager {
if request.source == HealRequestSource::Admin && request.force_start {
let overlapping: Vec<String> = {
let active_heals = self.active_heals.lock().await;
active_heals
let queue = self.heal_queue.lock().await;
let retrying = self.retrying_heals.lock().await;
let mut ids = active_heals
.iter()
.filter(|(task_id, task)| {
task.source == HealRequestSource::Admin
@@ -1604,7 +1670,17 @@ impl HealManager {
&& *task_id != &request.id
})
.map(|(task_id, _)| task_id.clone())
.collect()
.collect::<Vec<_>>();
ids.extend(
queue
.requests()
.chain(retrying.values().map(|retrying| &retrying.request))
.filter(|pending| {
root_recovery::is_root_heal(&pending.heal_type, pending.source) && pending.id != request.id
})
.map(|pending| pending.id.clone()),
);
ids
};
for task_id in overlapping {
match self.cancel_task(&task_id).await {
@@ -1618,17 +1694,14 @@ impl HealManager {
result = "force_start_cancelled_overlap",
"Admin forceStart cancelled an overlapping heal task"
),
Err(err) => warn!(
target: "rustfs::heal::manager",
event = EVENT_HEAL_QUEUE_ADMISSION,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_MANAGER,
request_id = %request.id,
cancelled_task_id = %task_id,
error = %err,
result = "force_start_cancel_failed",
"Admin forceStart failed to cancel an overlapping heal task"
),
Err(err) => return Err(err),
}
}
// A failed or timed-out replay may have only its durable owner
// left. Root responsibility overlaps every administrator path.
for pending in self.root_recovery.pending().await? {
if pending.id != request.id {
self.cancel_task(&pending.id).await?;
}
}
}
@@ -1641,6 +1714,9 @@ impl HealManager {
// active -> retrying transitions can slip between duplicate checks.
let lock_phase_start = Instant::now();
let active_heals = self.active_heals.lock().await;
if self.cancel_token.is_cancelled() {
return Err(Error::Other("Heal manager is stopping".to_string()));
}
#[cfg(test)]
pause_duplicate_admission_after_active_lock(&request.id).await;
let mut queue = self.heal_queue.lock().await;
@@ -2145,6 +2221,7 @@ impl HealManager {
{
let mut active_heals = self.active_heals.lock().await;
if let Some(task) = active_heals.get(&canonical_task_id) {
self.root_recovery.remove(&task.id, &task.heal_type, task.source).await?;
task.cancel().await?;
let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await;
publish_completed_heal(&self.completed_heals, &self.task_aliases, &canonical_task_id, completed, true).await;
@@ -2168,6 +2245,11 @@ impl HealManager {
{
let mut retrying_heals = self.retrying_heals.lock().await;
if let Some(retrying) = retrying_heals.get(&canonical_task_id) {
self.root_recovery
.remove(&canonical_task_id, &retrying.request.heal_type, retrying.request.source)
.await?;
}
if let Some(retrying) = retrying_heals.remove(&canonical_task_id) {
retrying.cancel_token.cancel();
drop(retrying_heals);
@@ -2188,6 +2270,11 @@ impl HealManager {
}
let mut queue = self.heal_queue.lock().await;
if let Some(request) = queue.requests().find(|request| request.id == canonical_task_id) {
self.root_recovery
.remove(&request.id, &request.heal_type, request.source)
.await?;
}
if queue.remove_request_id(&canonical_task_id).is_some() {
publish_heal_queue_length(&queue);
info!(
@@ -2205,6 +2292,10 @@ impl HealManager {
return Ok(());
}
drop(queue);
if self.root_recovery.cancel_pending(&canonical_task_id).await? {
return Ok(());
}
Err(Error::TaskNotFound {
task_id: task_id.to_string(),
})
@@ -2223,6 +2314,7 @@ impl HealManager {
for task_id in &task_ids {
if let Some(task) = active_heals.get(task_id) {
self.root_recovery.remove(&task.id, &task.heal_type, task.source).await?;
task.cancel().await?;
let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await;
publish_completed_heal(&self.completed_heals, &self.task_aliases, task_id, completed, true).await;
@@ -2251,6 +2343,11 @@ impl HealManager {
.collect::<Vec<_>>();
for task_id in &task_ids {
if let Some(retrying) = retrying_heals.get(task_id) {
self.root_recovery
.remove(task_id, &retrying.request.heal_type, retrying.request.source)
.await?;
}
if let Some(retrying) = retrying_heals.remove(task_id) {
retrying.cancel_token.cancel();
cancelled += 1;
@@ -2273,6 +2370,14 @@ impl HealManager {
}
let mut queue = self.heal_queue.lock().await;
for request in queue
.requests()
.filter(|request| heal_type_matches_path(&request.heal_type, heal_path))
{
self.root_recovery
.remove(&request.id, &request.heal_type, request.source)
.await?;
}
let queued_cancelled = queue.remove_matching(|request| heal_type_matches_path(&request.heal_type, heal_path));
if !queued_cancelled.is_empty() {
publish_heal_queue_length(&queue);
@@ -2284,6 +2389,13 @@ impl HealManager {
self.remove_mrf_repair_notice_targets_for_task(&request.id);
}
if heal_type_matches_path(&HealType::Cluster, heal_path) {
for pending in self.root_recovery.pending().await? {
if self.root_recovery.cancel_pending(&pending.id).await? {
cancelled += 1;
}
}
}
if cancelled == 0 {
return Err(Error::TaskNotFound {
task_id: heal_path.to_string(),
@@ -2395,6 +2507,7 @@ impl std::fmt::Debug for HealManager {
mod auto_scan;
mod queue;
mod root_recovery;
mod scheduler;
mod unclean_shutdown;
@@ -0,0 +1,342 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Graceful-shutdown handoff for administrator root heals. This namespace is
//! separate from erasure-set checkpoints and replacement generations, which
//! cannot represent a cluster traversal. One coordinator disk owns each
//! record; never create a fallback copy after an uncertain write or deletion.
use super::*;
use crate::heal::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes};
use crate::heal::{DiskStore, RUSTFS_META_BUCKET};
use serde::{Deserialize, Serialize};
// The metadata bucket already exists and its parent is durable. Creating a
// nested journal directory here would also require syncing every ancestor.
const ROOT_RECOVERY_PREFIX: &str = "root-heal-";
const ROOT_RECOVERY_SCHEMA: u32 = 1;
#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct RootHealIntent {
schema: u32,
task_id: String,
#[serde(deserialize_with = "decode_options")]
options: HealOptions,
priority: HealPriority,
retry_attempts: u32,
created_at: SystemTime,
}
impl RootHealIntent {
fn from_request(request: &HealRequest) -> Self {
Self {
schema: ROOT_RECOVERY_SCHEMA,
task_id: request.id.clone(),
options: request.options.clone(),
priority: request.priority,
retry_attempts: request.retry_attempts,
created_at: request.created_at,
}
}
fn into_request(self) -> HealRequest {
let mut request = HealRequest::new(HealType::Cluster, self.options, self.priority);
request.id = self.task_id;
request.source = HealRequestSource::Admin;
request.retry_attempts = self.retry_attempts;
request.created_at = self.created_at;
request
}
}
#[derive(Default)]
pub(super) struct RootHealRecovery {
mutation: Mutex<()>,
#[cfg(test)]
disks: Option<Vec<DiskStore>>,
}
pub(super) fn is_root_heal(heal_type: &HealType, source: HealRequestSource) -> bool {
source == HealRequestSource::Admin && matches!(heal_type, HealType::Cluster)
}
fn decode_options<'de, D: serde::Deserializer<'de>>(deserializer: D) -> std::result::Result<HealOptions, D::Error> {
let value = serde_json::Value::deserialize(deserializer)?;
let object = value
.as_object()
.ok_or_else(|| serde::de::Error::custom("root heal options must be an object"))?;
const FIELDS: &[&str] = &[
"scan_mode",
"remove_corrupted",
"recreate_missing",
"update_parity",
"recursive",
"dry_run",
"no_lock",
"timeout",
"pool_index",
"set_index",
];
if object.keys().any(|key| !FIELDS.contains(&key.as_str())) {
return Err(serde::de::Error::custom("unknown root heal recovery option"));
}
let options: HealOptions = serde_json::from_value(value).map_err(serde::de::Error::custom)?;
if options.no_lock {
return Err(serde::de::Error::custom("administrator root heal cannot skip namespace locking"));
}
Ok(options)
}
fn intent_path(task_id: &str) -> Result<String> {
let parsed = uuid::Uuid::parse_str(task_id).map_err(|_| Error::Other("Invalid root heal recovery task id".to_string()))?;
if parsed.to_string() != task_id {
return Err(Error::Other("Noncanonical root heal recovery task id".to_string()));
}
Ok(format!("{ROOT_RECOVERY_PREFIX}{task_id}.json"))
}
fn decode_intent(task_id: &str, bytes: &[u8]) -> Result<RootHealIntent> {
let _ = intent_path(task_id)?;
let intent: RootHealIntent = serde_json::from_slice(bytes)
.map_err(|error| Error::Other(format!("Invalid root heal recovery record {task_id}: {error}")))?;
if intent.schema != ROOT_RECOVERY_SCHEMA || intent.task_id != task_id {
return Err(Error::Other(format!("Unsupported or mismatched root heal recovery record {task_id}")));
}
Ok(intent)
}
impl RootHealRecovery {
#[cfg(test)]
pub(super) fn with_disks(disks: Vec<DiskStore>) -> Self {
Self {
mutation: Mutex::new(()),
disks: Some(disks),
}
}
async fn disks(&self) -> Result<Vec<DiskStore>> {
#[cfg(test)]
if let Some(disks) = &self.disks {
return Ok(disks.clone());
}
let map = local_disk_map_read().await;
if map.values().any(Option::is_none) {
return Err(Error::Other("Root heal recovery owner may be on an unavailable local disk".to_string()));
}
let mut disks = map.values().flatten().cloned().collect::<Vec<_>>();
disks.sort_by_key(|disk| EcstoreDiskAPI::endpoint(disk.as_ref()).to_string());
Ok(disks)
}
async fn find(disks: &[DiskStore], task_id: &str) -> Result<Option<(DiskStore, EcstoreDiskBytes)>> {
let path = intent_path(task_id)?;
let mut found = None;
for disk in disks {
// read_all reports FileNotFound even when the whole metadata
// volume is absent; that is an unknown owner, not empty state.
EcstoreDiskAPI::stat_volume(disk.as_ref(), RUSTFS_META_BUCKET).await?;
match EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, &path).await {
Ok(bytes) => {
decode_intent(task_id, &bytes)?;
if found.is_some() {
return Err(Error::Other(format!("Multiple root heal recovery owners for {task_id}")));
}
found = Some((disk.clone(), bytes));
}
Err(DiskError::FileNotFound) => {}
Err(error) => return Err(Error::Disk(error)),
}
}
Ok(found)
}
pub(super) async fn persist(&self, request: &HealRequest) -> Result<()> {
if !is_root_heal(&request.heal_type, request.source) {
return Ok(());
}
let _guard = self.mutation.lock().await;
let disks = self.disks().await?;
let existing = Self::find(&disks, &request.id).await?;
let (disk, expected) = match existing {
Some((disk, bytes)) => (disk, Some(bytes)),
None => {
let disk = disks
.first()
.cloned()
.ok_or_else(|| Error::Other("No local disk available for root heal shutdown recovery".to_string()))?;
(disk, None)
}
};
if request.options.no_lock {
return Err(Error::Other("Administrator root heal cannot skip namespace locking".to_string()));
}
let bytes = serde_json::to_vec(&RootHealIntent::from_request(request))
.map_err(|error| Error::Other(format!("Serialize root heal recovery record: {error}")))?;
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&intent_path(&request.id)?,
expected,
Some(bytes.into()),
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(()),
_ => Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
}
}
pub(super) async fn remove(&self, task_id: &str, heal_type: &HealType, source: HealRequestSource) -> Result<bool> {
if !is_root_heal(heal_type, source) {
return Ok(false);
}
let _guard = self.mutation.lock().await;
let Some((disk, bytes)) = Self::find(&self.disks().await?, task_id).await? else {
return Ok(false);
};
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&intent_path(task_id)?,
Some(bytes),
None,
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(true),
_ => Err(Error::Other(format!("Root heal recovery record changed while retiring {task_id}"))),
}
}
pub(super) async fn checkpoint_failed_execution(&self, task: &HealTask) -> Result<()> {
if !is_root_heal(&task.heal_type, task.source) {
return Ok(());
}
let remaining = match task.retry_request_with_remaining_timeout().await {
Ok(request) => request.options.timeout,
Err(Error::TaskTimeout) => Some(Duration::ZERO),
Err(error) => return Err(error),
};
let _guard = self.mutation.lock().await;
let Some((disk, expected)) = Self::find(&self.disks().await?, &task.id).await? else {
// A first execution that failed has no restart handoff to update.
return Ok(());
};
let mut intent = decode_intent(&task.id, &expected)?;
let mut expected_options = intent.options.clone();
expected_options.timeout = task.options.timeout;
if intent.created_at != task.created_at || intent.priority != task.priority || expected_options != task.options {
return Err(Error::Other(format!("Root heal recovery owner changed for {}", task.id)));
}
// A terminal timeout leaves no runtime owner for stop() to snapshot.
// Checkpoint its consumed budget before publishing terminal status;
// never refund time if an earlier checkpoint is already stricter.
intent.options.timeout = match (intent.options.timeout, remaining) {
(Some(previous), Some(remaining)) => Some(previous.min(remaining)),
(previous, remaining) => previous.or(remaining),
};
intent.retry_attempts = intent.retry_attempts.max(task.retry_attempts);
let bytes = serde_json::to_vec(&intent)
.map_err(|error| Error::Other(format!("Serialize root heal recovery checkpoint: {error}")))?;
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&intent_path(&task.id)?,
Some(expected),
Some(bytes.into()),
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(()),
_ => Err(Error::Other(format!("Root heal recovery record changed while checkpointing {}", task.id))),
}
}
pub(super) async fn cancel_pending(&self, task_id: &str) -> Result<bool> {
if intent_path(task_id).is_err() {
return Ok(false);
}
self.remove(task_id, &HealType::Cluster, HealRequestSource::Admin).await
}
pub(super) async fn pending(&self) -> Result<Vec<HealRequest>> {
let _guard = self.mutation.lock().await;
let disks = self.disks().await?;
let mut ids = HashSet::new();
for disk in &disks {
EcstoreDiskAPI::stat_volume(disk.as_ref(), RUSTFS_META_BUCKET).await?;
let entries = match EcstoreDiskAPI::list_dir(disk.as_ref(), "", RUSTFS_META_BUCKET, "", -1).await {
Ok(entries) => entries,
Err(DiskError::FileNotFound) => continue,
Err(error) => return Err(Error::Disk(error)),
};
for entry in entries {
let Some(task_id) = entry
.strip_prefix(ROOT_RECOVERY_PREFIX)
.and_then(|entry| entry.strip_suffix(".json"))
else {
continue;
};
let _ = intent_path(task_id)?;
ids.insert(task_id.to_string());
}
}
let mut requests = Vec::new();
for task_id in ids {
if let Some((_, bytes)) = Self::find(&disks, &task_id).await? {
requests.push(decode_intent(&task_id, &bytes)?.into_request());
}
}
requests.sort_by(|left, right| left.created_at.cmp(&right.created_at).then_with(|| left.id.cmp(&right.id)));
Ok(requests)
}
}
impl HealManager {
pub(super) async fn replay_root_heals(&self) -> Result<()> {
// Decode every record before admitting anything. These are already
// accepted responsibilities, so restore distinct IDs even when their
// paths overlap or the configured admission capacity has changed.
let requests = self.root_recovery.pending().await?;
let active = self.active_heals.lock().await;
let mut queue = self.heal_queue.lock().await;
let retrying = self.retrying_heals.lock().await;
for mut request in requests {
request.force_start = true;
let existing = active
.get(&request.id)
.map(|task| request_matches_task(&request, task))
.or_else(|| {
queue
.requests()
.find(|queued| queued.id == request.id)
.map(|queued| request_matches_request(&request, queued))
})
.or_else(|| {
retrying
.get(&request.id)
.map(|retrying| request_matches_request(&request, &retrying.request))
});
match existing {
Some(true) => continue,
Some(false) => return Err(Error::Other(format!("Conflicting root heal recovery task {}", request.id))),
None => {}
}
queue.push(request);
}
publish_heal_queue_length(&queue);
Ok(())
}
}
+40
View File
@@ -26,6 +26,7 @@ impl HealManager {
let retrying_heals = self.retrying_heals.clone();
let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone();
let replacement_recovery_anchors = self.replacement_recovery_anchors.clone();
let root_recovery = self.root_recovery.clone();
let cancel_token = self.cancel_token.clone();
let statistics = self.statistics.clone();
let storage = self.storage.clone();
@@ -59,6 +60,7 @@ impl HealManager {
retrying_heals: &retrying_heals,
mrf_repair_notice_targets: &mrf_repair_notice_targets,
replacement_recovery_anchors: &replacement_recovery_anchors,
root_recovery: &root_recovery,
config: &config,
statistics: &statistics,
storage: &storage,
@@ -78,6 +80,7 @@ impl HealManager {
retrying_heals: &retrying_heals,
mrf_repair_notice_targets: &mrf_repair_notice_targets,
replacement_recovery_anchors: &replacement_recovery_anchors,
root_recovery: &root_recovery,
config: &config,
statistics: &statistics,
storage: &storage,
@@ -106,6 +109,7 @@ impl HealManager {
retrying_heals,
mrf_repair_notice_targets,
replacement_recovery_anchors,
root_recovery,
config,
statistics,
storage,
@@ -117,6 +121,9 @@ impl HealManager {
let config = config.read().await;
let mainline_pressure = Self::mainline_throttle_active(&config, workload_provider);
let mut active_heals_guard = active_heals.lock().await;
if cancel_token.is_cancelled() {
return;
}
publish_active_heal_count(&active_heals_guard);
// Check if new heal tasks can be started
@@ -206,6 +213,7 @@ impl HealManager {
let replacement_recovery_anchors_clone = replacement_recovery_anchors.clone();
let statistics_clone = statistics.clone();
let notify_clone = notify.clone();
let root_recovery_clone = root_recovery.clone();
let manager_cancel_token = cancel_token.clone();
let task_type_label_for_spawn = task_type_label.clone();
let task_set_label_for_spawn = task_set_label.clone();
@@ -294,6 +302,38 @@ impl HealManager {
tests::pause_completed_retention_before_publish(&task_id, &completed_status).await;
let mut active_heals_guard = active_heals_clone.lock().await;
let owns_completion = active_heals_guard.contains_key(&task_id);
if owns_completion
&& result.is_ok()
&& let Err(error) = root_recovery_clone.remove(&task_id, &task.heal_type, task.source).await
{
// Keep the durable responsibility if retirement fails.
// Replaying a completed traversal is idempotent.
warn!(
target: "rustfs::heal::manager",
event = EVENT_HEAL_SCHEDULER_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_MANAGER,
task_id,
state = "root_recovery_retirement_failed",
error = %error,
"Failed to retire root heal recovery record"
);
}
if owns_completion
&& result.is_err()
&& let Err(error) = root_recovery_clone.checkpoint_failed_execution(&task).await
{
warn!(
target: "rustfs::heal::manager",
event = EVENT_HEAL_SCHEDULER_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_MANAGER,
task_id,
state = "root_recovery_checkpoint_failed",
error = %error,
"Failed to checkpoint root heal recovery execution budget"
);
}
let cancelled_completion = if owns_completion {
false
} else {
+2
View File
@@ -26,6 +26,7 @@ use rustfs_madmin::heal_commands::HealResultItem;
use std::sync::Mutex as StdMutex;
use tempfile::TempDir;
mod root_recovery;
mod running_mainline;
use super::super::{DiskOption, DiskStore, Endpoint, new_disk, storage_api::status::BucketInfo};
@@ -94,6 +95,7 @@ async fn process_manager_queue_once(manager: &HealManager) {
retrying_heals: &manager.retrying_heals,
mrf_repair_notice_targets: &manager.mrf_repair_notice_targets,
replacement_recovery_anchors: &manager.replacement_recovery_anchors,
root_recovery: &manager.root_recovery,
config: &manager.config,
statistics: &manager.statistics,
storage: &manager.storage,
@@ -0,0 +1,465 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::super::root_recovery::RootHealRecovery;
use super::*;
use crate::heal::RUSTFS_META_BUCKET;
async fn recovery_disk() -> (TempDir, DiskStore) {
let temp = TempDir::new().expect("temporary root recovery disk");
let endpoint = Endpoint::try_from(temp.path().to_string_lossy().as_ref()).expect("disk endpoint");
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("local recovery disk");
match disk.make_volume(RUSTFS_META_BUCKET).await {
Ok(()) | Err(DiskError::VolumeExists) => {}
Err(error) => panic!("metadata volume: {error}"),
}
(temp, disk)
}
fn recovery_manager(disks: Vec<DiskStore>) -> HealManager {
let mut manager = HealManager::new(
Arc::new(MockStorage),
Some(HealConfig {
enable_auto_heal: false,
..Default::default()
}),
);
manager.root_recovery = Arc::new(RootHealRecovery::with_disks(disks));
manager
}
fn root_request() -> HealRequest {
let mut request = HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::High);
request.source = HealRequestSource::Admin;
request
}
async fn active_root(manager: &HealManager, request: HealRequest) -> Arc<HealTask> {
let task = Arc::new(HealTask::from_request(request, manager.storage.clone()));
*task.status.write().await = HealTaskStatus::Running;
task.progress.write().await.update_object_progress(1, 1, 0, 0, 128);
manager.active_heals.lock().await.insert(task.id.clone(), task.clone());
task
}
#[tokio::test]
async fn root_recovery_shutdown_restart_replays_same_id_and_success_retires_intent() {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let mut request = root_request();
request.options.recursive = true;
let task = active_root(&manager, request.clone()).await;
manager.stop().await.expect("durable shutdown handoff");
assert!(task.cancel_token.is_cancelled());
drop(manager);
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("replay durable root");
restarted.replay_root_heals().await.expect("replay is idempotent");
assert_eq!(restarted.get_queue_length().await, 1);
let restored = restarted
.heal_queue
.lock()
.await
.requests()
.next()
.cloned()
.expect("restored request");
assert_eq!(restored.id, request.id);
assert_eq!(restored.options, request.options);
assert_eq!(restored.priority, request.priority);
assert_eq!(restored.retry_attempts, request.retry_attempts);
assert_eq!(restored.created_at, request.created_at);
process_manager_queue_once(&restarted).await;
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if matches!(restarted.get_task_status(&request.id).await, Ok(HealTaskStatus::Completed))
&& !restarted.active_heals.lock().await.contains_key(&request.id)
{
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("restored root executes successfully");
assert!(restarted.root_recovery.pending().await.expect("read completion").is_empty());
}
#[tokio::test]
async fn root_recovery_explicit_cancel_covers_active_queued_retrying_and_durable_only() {
for state in ["active", "queued", "retrying", "durable_only", "root_path"] {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let request = root_request();
manager.root_recovery.persist(&request).await.expect("durable responsibility");
match state {
"active" => {
active_root(&manager, request.clone()).await;
}
"queued" => {
manager.replay_root_heals().await.expect("queued recovery");
}
"retrying" => {
insert_retrying_request(&manager, request.clone()).await;
}
_ => {}
}
if state == "root_path" {
assert_eq!(manager.cancel_tasks_for_path("").await.expect("cancel durable root path"), 1);
} else {
manager.cancel_task(&request.id).await.expect("cancel root responsibility");
}
drop(manager);
let restarted = recovery_manager(vec![disk]);
restarted
.replay_root_heals()
.await
.expect("restart after explicit cancellation");
assert_eq!(restarted.get_queue_length().await, 0, "state={state}");
}
}
#[tokio::test]
async fn root_recovery_force_start_cancels_durable_only_responsibility() {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let old = root_request();
manager
.root_recovery
.persist(&old)
.await
.expect("old terminal responsibility");
let mut new = root_request();
new.force_start = true;
assert_eq!(
manager
.submit_heal_request(new.clone())
.await
.expect("force start replacement"),
HealAdmissionResult::Accepted
);
assert!(manager.root_recovery.pending().await.expect("old owner retired").is_empty());
manager.stop().await.expect("persist new root only");
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("restart replacement");
let ids = restarted
.heal_queue
.lock()
.await
.requests()
.map(|request| request.id.clone())
.collect::<Vec<_>>();
assert_eq!(ids, [new.id]);
}
#[tokio::test]
async fn root_recovery_force_start_replaces_fresh_queued_and_retrying_admin_roots() {
for retrying in [false, true] {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let old = root_request();
if retrying {
insert_retrying_request(&manager, old.clone()).await;
} else {
manager.submit_heal_request(old.clone()).await.expect("queue original root");
}
assert!(manager.root_recovery.pending().await.expect("not handed off yet").is_empty());
let mut new = root_request();
new.force_start = true;
assert_eq!(
manager.submit_heal_request(new.clone()).await.expect("force replacement"),
HealAdmissionResult::Accepted
);
manager.stop().await.expect("handoff only the new responsibility");
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("restart after forceStart");
let ids = restarted
.heal_queue
.lock()
.await
.requests()
.map(|request| request.id.clone())
.collect::<Vec<_>>();
assert_eq!(ids, [new.id], "retrying={retrying}; old={}", old.id);
}
}
#[tokio::test]
async fn root_recovery_invalid_records_are_retained_without_partial_replay() {
for kind in ["truncated", "schema", "identity", "option", "no_lock"] {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let valid = root_request();
let invalid = root_request();
manager.root_recovery.persist(&valid).await.expect("valid root record");
manager
.root_recovery
.persist(&invalid)
.await
.expect("record before corruption");
let path = format!("root-heal-{}.json", invalid.id);
let original = disk.read_all(RUSTFS_META_BUCKET, &path).await.expect("read root record");
let mut value: serde_json::Value = serde_json::from_slice(&original).expect("record JSON");
match kind {
"schema" => value["schema"] = 2.into(),
"identity" => value["task_id"] = valid.id.clone().into(),
"option" => value["options"]["future_delete_mode"] = true.into(),
"no_lock" => value["options"]["no_lock"] = true.into(),
_ => {}
}
let bytes = if kind == "truncated" {
b"{".to_vec()
} else {
serde_json::to_vec(&value).expect("modified record")
};
disk.write_all(RUSTFS_META_BUCKET, &path, bytes.clone().into())
.await
.expect("inject bad record");
assert!(manager.replay_root_heals().await.is_err(), "kind={kind}");
assert_eq!(manager.get_queue_length().await, 0, "no partial admission for {kind}");
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("bad record retained")
.as_ref(),
bytes
);
let mut forced = root_request();
forced.force_start = true;
assert!(
manager.submit_heal_request(forced).await.is_err(),
"forceStart must not discard unknown state"
);
}
}
#[tokio::test]
async fn root_recovery_failed_handoff_keeps_runtime_owner_and_does_not_try_another_disk() {
let (_temp, disk) = recovery_disk().await;
let (unavailable_temp, unavailable) = recovery_disk().await;
std::fs::remove_dir_all(unavailable_temp.path().join(RUSTFS_META_BUCKET)).expect("make owner volume unavailable");
let manager = recovery_manager(vec![unavailable, disk.clone()]);
let task = active_root(&manager, root_request()).await;
assert!(manager.stop().await.is_err());
assert!(
manager.cancel_task(&task.id).await.is_err(),
"missing owner cannot acknowledge cancellation"
);
assert!(!manager.cancel_token.is_cancelled());
assert!(!task.cancel_token.is_cancelled());
assert!(manager.active_heals.lock().await.contains_key(&task.id));
assert!(
RootHealRecovery::with_disks(vec![disk])
.pending()
.await
.expect("other disk remains empty")
.is_empty()
);
}
#[tokio::test]
async fn root_recovery_shutdown_fences_new_admission_and_preserves_later_cancellation() {
for operation_kind in ["submit", "force_start", "cancel"] {
let cancel = operation_kind == "cancel";
let (_temp, disk) = recovery_disk().await;
let manager = Arc::new(recovery_manager(vec![disk.clone()]));
let request = root_request();
active_root(&manager, request.clone()).await;
let queue = manager.heal_queue.lock().await;
let stopping = manager.clone();
let stop = tokio::spawn(async move { stopping.stop().await });
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if manager.active_heals.try_lock().is_err() {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("shutdown owns active lock while waiting for queue");
let concurrent = manager.clone();
let operation = tokio::spawn(async move {
if cancel {
concurrent.cancel_task(&request.id).await
} else {
let mut new = root_request();
new.force_start = operation_kind == "force_start";
concurrent.submit_heal_request(new).await.map(|_| ())
}
});
drop(queue);
stop.await.expect("shutdown task").expect("durable shutdown");
let result = operation.await.expect("concurrent operation");
assert_eq!(result.is_ok(), cancel, "operation={operation_kind}");
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("read final responsibility");
assert_eq!(restarted.get_queue_length().await, usize::from(!cancel));
}
}
#[tokio::test]
async fn root_recovery_exhausted_timeout_is_not_reset_by_restart() {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let mut request = root_request();
request.options.timeout = Some(Duration::from_secs(10));
let task = active_root(&manager, request.clone()).await;
task.set_execution_elapsed_for_test(Duration::from_secs(11)).await;
manager.stop().await.expect("persist exhausted execution budget");
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("restore bounded request");
assert_eq!(
restarted
.heal_queue
.lock()
.await
.requests()
.next()
.expect("restored root")
.options
.timeout,
Some(Duration::ZERO)
);
process_manager_queue_once(&restarted).await;
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if matches!(restarted.get_task_status(&request.id).await, Ok(HealTaskStatus::Timeout)) {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("exhausted request stays timed out");
restarted
.cancel_task(&request.id)
.await
.expect("timeout responsibility remains cancellable");
assert!(restarted.root_recovery.pending().await.expect("retired timeout").is_empty());
}
#[tokio::test]
async fn root_recovery_shutdown_preserves_remaining_execution_budget() {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let mut request = root_request();
request.options.timeout = Some(Duration::from_secs(60));
let task = active_root(&manager, request).await;
task.set_execution_elapsed_for_test(Duration::from_secs(20)).await;
manager.stop().await.expect("handoff with consumed execution time");
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("restore remaining budget");
let queue = restarted.heal_queue.lock().await;
let remaining = queue
.requests()
.next()
.expect("restored root")
.options
.timeout
.expect("remaining timeout");
assert!(remaining <= Duration::from_secs(40), "elapsed execution must not be refunded");
assert!(
remaining >= Duration::from_secs(30),
"shutdown fixture should retain most of its remaining budget"
);
}
#[tokio::test]
async fn root_recovery_force_start_after_shutdown_does_not_retire_original_owner() {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let old = root_request();
active_root(&manager, old.clone()).await;
manager.stop().await.expect("handoff original root");
let mut new = root_request();
new.force_start = true;
assert!(manager.submit_heal_request(new).await.is_err());
let restarted = recovery_manager(vec![disk]);
restarted.replay_root_heals().await.expect("original responsibility remains");
let ids = restarted
.heal_queue
.lock()
.await
.requests()
.map(|request| request.id.clone())
.collect::<Vec<_>>();
assert_eq!(ids, [old.id]);
}
#[tokio::test]
async fn root_recovery_terminal_timeout_updates_only_existing_journal_before_second_restart() {
for durable in [false, true] {
let (_temp, disk) = recovery_disk().await;
let manager = recovery_manager(vec![disk.clone()]);
let mut request = root_request();
request.options.timeout = Some(Duration::from_nanos(1));
if durable {
manager
.root_recovery
.persist(&request)
.await
.expect("persist nonzero execution budget");
manager.replay_root_heals().await.expect("first restart");
} else {
manager
.submit_heal_request(request.clone())
.await
.expect("first root execution");
}
process_manager_queue_once(&manager).await;
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if matches!(manager.get_task_status(&request.id).await, Ok(HealTaskStatus::Timeout))
&& !manager.active_heals.lock().await.contains_key(&request.id)
{
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("real execution exhausts a nonzero budget");
assert!(!manager.active_heals.lock().await.contains_key(&request.id));
let restarted = recovery_manager(vec![disk]);
restarted
.replay_root_heals()
.await
.expect("second restart after terminal timeout");
let queue = restarted.heal_queue.lock().await;
if durable {
assert_eq!(
queue
.requests()
.next()
.expect("remaining timeout responsibility")
.options
.timeout,
Some(Duration::ZERO)
);
} else {
assert!(queue.is_empty(), "terminal failure must not create a new durable responsibility");
}
}
}
+5
View File
@@ -513,6 +513,11 @@ impl HealTask {
}
}
#[cfg(test)]
pub(crate) async fn set_execution_elapsed_for_test(&self, elapsed: Duration) {
*self.task_start_instant.write().await = Some(Instant::now() - elapsed);
}
pub(crate) async fn retry_request_with_remaining_timeout(&self) -> Result<HealRequest> {
let mut request = self.retry_request();
if self.options.timeout.is_some() {
+5 -1
View File
@@ -61,10 +61,14 @@ pub fn create_ahm_services_cancel_token() -> CancellationToken {
}
/// Shutdown all heal services gracefully
pub fn shutdown_ahm_services() {
pub async fn shutdown_ahm_services() -> Result<()> {
if let Some(manager) = get_heal_manager() {
manager.stop().await?;
}
if let Some(cancel_token) = GLOBAL_AHM_SERVICES_CANCEL_TOKEN.get() {
cancel_token.cancel();
}
Ok(())
}
struct HealRuntime {
+432
View File
@@ -0,0 +1,432 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use bytes::Bytes;
use datafusion::object_store::{Error, Result};
use futures::{Stream, StreamExt, stream::BoxStream};
use transform_stream::AsyncTryStream;
use crate::SelectError;
/// Arrow accepts byte-sized CSV controls. Unicode quotes need streaming normalization.
pub fn csv_input_requires_normalization(quote: Option<&str>, escape: Option<&str>) -> bool {
quote.is_some_and(|quote| quote.len() > 1) || escape.is_some_and(|escape| escape.len() > 1)
}
/// CSV syntax independent of request headers, serialization formats, or S3 DTOs.
#[derive(Default)]
pub(crate) struct CsvSyntax<'a> {
pub quote: Option<&'a str>,
pub escape: Option<&'a str>,
pub field: Option<&'a str>,
pub record: Option<&'a str>,
pub comment: Option<u8>,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum State {
FieldStart,
Unquoted,
Quoted,
AfterQuote,
Escaped,
Comment,
}
/// Emits ordinary CSV with every field quoted. This avoids reserving a sentinel
/// byte that might also appear in a UTF-8 field. Only a partial control token is
/// retained between chunks; neither records nor objects are buffered.
struct CsvInputNormalizer {
quote: Vec<u8>,
escape: Vec<u8>,
field: Vec<u8>,
record: Vec<u8>,
comment: Option<u8>,
default_records: bool,
state: State,
record_start: bool,
carry: Vec<u8>,
token_size: usize,
}
impl CsvInputNormalizer {
fn new(csv: &CsvSyntax<'_>) -> Self {
let quote = csv
.quote
.filter(|value| !value.is_empty())
.unwrap_or("\"")
.as_bytes()
.to_vec();
let escape = csv
.escape
.filter(|value| !value.is_empty())
.unwrap_or("\"")
.as_bytes()
.to_vec();
let field = csv.field.filter(|value| !value.is_empty()).unwrap_or(",").as_bytes().to_vec();
let record = csv
.record
.filter(|value| !value.is_empty())
.unwrap_or("\n")
.as_bytes()
.to_vec();
let token_size = quote.len().max(escape.len()).max(field.len()).max(record.len()).max(2);
Self {
quote,
escape,
field,
record,
comment: csv.comment,
default_records: csv.record.is_none(),
state: State::FieldStart,
record_start: true,
carry: Vec::new(),
token_size,
}
}
fn record_len(&self, bytes: &[u8]) -> usize {
if self.default_records && bytes.starts_with(b"\r\n") {
2
} else if self.default_records && bytes.starts_with(b"\r") {
1
} else if bytes.starts_with(&self.record) {
self.record.len()
} else {
0
}
}
fn push_value(output: &mut Vec<u8>, bytes: &[u8]) {
for byte in bytes {
if *byte == b'"' {
output.push(b'"');
}
output.push(*byte);
}
}
fn convert(&mut self, chunk: &[u8], last: bool) -> std::result::Result<Vec<u8>, SelectError> {
let mut bytes = std::mem::take(&mut self.carry);
bytes.extend_from_slice(chunk);
let end = if last {
bytes.len()
} else {
bytes.len().saturating_sub(self.token_size - 1)
};
let mut output = Vec::with_capacity(bytes.len());
let mut pos = 0;
while pos < end {
let rest = &bytes[pos..];
let record_len = self.record_len(rest);
let field = rest.starts_with(&self.field) && self.field.len() > record_len;
match self.state {
State::Comment => {
if record_len > 0 {
self.state = State::FieldStart;
pos += record_len;
} else {
pos += 1;
}
}
State::Escaped => {
if record_len > 0 {
return Err(SelectError::CsvParsingError);
}
Self::push_value(&mut output, &rest[..1]);
self.state = State::Quoted;
pos += 1;
}
State::Quoted if rest.starts_with(&self.quote) => {
self.state = State::AfterQuote;
pos += self.quote.len();
}
State::Quoted if rest.starts_with(&self.escape) => {
self.state = State::Escaped;
pos += self.escape.len();
}
State::Quoted => {
if record_len > 0 {
return Err(SelectError::CsvParsingError);
}
Self::push_value(&mut output, &rest[..1]);
pos += 1;
}
State::AfterQuote if rest.starts_with(&self.quote) => {
Self::push_value(&mut output, &self.quote);
self.state = State::Quoted;
pos += self.quote.len();
}
State::FieldStart if self.record_start && self.comment == Some(rest[0]) => {
self.state = State::Comment;
pos += 1;
}
State::FieldStart if rest.starts_with(&self.quote) => {
output.push(b'"');
self.state = State::Quoted;
self.record_start = false;
pos += self.quote.len();
}
_ if field || record_len > 0 => {
if self.state == State::FieldStart {
if field || !self.record_start {
output.extend_from_slice(b"\"\"");
}
} else {
output.push(b'"');
}
output.push(if field { b',' } else { b'\n' });
self.state = State::FieldStart;
self.record_start = !field;
pos += if field { self.field.len() } else { record_len };
}
_ => {
if self.state == State::FieldStart {
output.push(b'"');
}
self.state = State::Unquoted;
self.record_start = false;
Self::push_value(&mut output, &rest[..1]);
pos += 1;
}
}
}
self.carry.extend_from_slice(&bytes[pos..]);
if last {
match self.state {
State::Quoted | State::Escaped => return Err(SelectError::CsvParsingError),
State::Unquoted | State::AfterQuote => output.push(b'"'),
State::FieldStart if !self.record_start => output.extend_from_slice(b"\"\""),
State::FieldStart | State::Comment => {}
}
}
Ok(output)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn normalize_chunks(csv: &CsvSyntax<'_>, input: &[u8], chunk_size: usize) -> Vec<u8> {
let mut normalizer = CsvInputNormalizer::new(csv);
let mut output = Vec::new();
for chunk in input.chunks(chunk_size) {
output.extend(normalizer.convert(chunk, false).expect("normalize complete CSV input"));
assert!(normalizer.carry.len() < normalizer.token_size, "only a partial token may be retained");
}
output.extend(normalizer.convert(&[], true).expect("finish complete CSV input"));
output
}
#[test]
fn unicode_csv_quotes_preserve_values_at_every_chunk_boundary() {
let cases = [
("ع", "\"", "عcol1ع,عcol2ع,عcol3ع\n", "\"col1\",\"col2\",\"col3\"\n"),
("ع", "\"", "\"left\",tail\n", "\"\"\"left\"\"\",\"tail\"\n"),
("ع", "\"", "عA,Bع,plain\n", "\"A,B\",\"plain\"\n"),
("ع", "\"", "عAععBع,tail\n", "\"AعB\",\"tail\"\n"),
("ع", "\"", "عA\"عBع,tail\n", "\"AعB\",\"tail\"\n"),
("ع", "\\", "عA\\\"Bع,tail\n", "\"A\"\"B\",\"tail\"\n"),
("\"", "", "\"A界\"B\",\"C\"\n", "\"A\"\"B\",\"C\"\n"),
("🦀", "🦀", "🦀A🦀🦀B🦀,C\n", "\"A🦀B\",\"C\"\n"),
("ع", "\"", "a\0b,عc\0\n", "\"a\0b\",\"c\0d\"\n"),
("ع", "\"", "AعB,tail\n", "\"AعB\",\"tail\"\n"),
("ع", "\"", "عaعsuffix,tail\n", "\"asuffix\",\"tail\"\n"),
("ع", "\"", ",\n", "\"\",\"\"\n"),
("ع", "\"", "a,", "\"a\",\"\""),
("ع", "\"", "عع", "\"\""),
("ع", "\"", "\n", "\n"),
("ع", "\"", "", ""),
];
for (quote, escape, input, expected) in cases {
let csv = CsvSyntax {
quote: Some(quote),
escape: Some(escape),
record: Some("\n"),
..Default::default()
};
for chunk_size in 1..=input.len().max(1) {
assert_eq!(
normalize_chunks(&csv, input.as_bytes(), chunk_size),
expected.as_bytes(),
"input={input:?}, chunk_size={chunk_size}"
);
}
}
}
#[test]
fn unicode_csv_quotes_keep_custom_delimiters_and_comments_out_of_values() {
let csv = CsvSyntax {
quote: Some("ع"),
escape: Some("\\"),
field: Some(""),
record: Some("^Y"),
comment: Some(b'#'),
};
let input = "#skipع界^Yعa界bع界\"literal\"^Yعline\nbreakع界end^Y";
let expected = "\"a界b\",\"\"\"literal\"\"\"\n\"line\nbreak\",\"end\"\n";
for chunk_size in 1..=input.len() {
assert_eq!(normalize_chunks(&csv, input.as_bytes(), chunk_size), expected.as_bytes());
}
}
#[test]
fn unicode_csv_quotes_reject_unterminated_fields_and_quoted_record_delimiters() {
for input in ["عunfinished", "عescape\\", "عline\nbreakع\n", "عline\\\nbreakع\n"] {
let csv = CsvSyntax {
quote: Some("ع"),
escape: Some("\\"),
..Default::default()
};
let mut normalizer = CsvInputNormalizer::new(&csv);
assert_eq!(normalizer.convert(input.as_bytes(), true), Err(SelectError::CsvParsingError));
}
}
#[test]
fn unicode_csv_quotes_preserve_omitted_syntax_defaults() {
assert!(!csv_input_requires_normalization(None, None));
assert!(!csv_input_requires_normalization(Some("\""), Some("\\")));
assert!(csv_input_requires_normalization(Some("ع"), None));
assert!(csv_input_requires_normalization(None, Some("")));
let quote_only = CsvSyntax {
quote: Some("ع"),
..Default::default()
};
assert_eq!(
normalize_chunks(&quote_only, "عA\"عBع,tail\r\n".as_bytes(), 1),
"\"AعB\",\"tail\"\n".as_bytes()
);
let escape_only = CsvSyntax {
escape: Some(""),
..Default::default()
};
assert_eq!(
normalize_chunks(&escape_only, "\"A界\"B\",tail\r\n".as_bytes(), 1),
b"\"A\"\"B\",\"tail\"\n"
);
}
#[test]
fn unicode_csv_quotes_stream_large_fields_without_retaining_records() {
let csv = CsvSyntax {
quote: Some("ع"),
..Default::default()
};
let mut normalizer = CsvInputNormalizer::new(&csv);
let chunk = vec![b'x'; 64 * 1024];
let mut output_len = normalizer.convert("ع".as_bytes(), false).expect("opening quote").len();
for _ in 0..64 {
let output = normalizer.convert(&chunk, false).expect("stream field chunk");
assert!(output.len() >= chunk.len() - 3, "field data must be emitted before its closing quote");
assert!(normalizer.carry.len() < 4);
output_len += output.len();
}
output_len += normalizer.convert("ع\n".as_bytes(), true).expect("close field").len();
assert_eq!(output_len, chunk.len() * 64 + 3);
}
}
pub(crate) fn normalize_csv_stream<S>(stream: S, csv: &CsvSyntax<'_>) -> BoxStream<'static, Result<Bytes>>
where
S: Stream<Item = Result<Bytes>> + Send + 'static,
{
let mut normalizer = CsvInputNormalizer::new(csv);
AsyncTryStream::<Bytes, Error, _>::new(|mut y| async move {
futures::pin_mut!(stream);
while let Some(chunk) = stream.next().await {
let converted = normalizer.convert(&chunk?, false).map_err(|source| Error::Generic {
store: "EcObjectStore",
source: Box::new(source),
})?;
if !converted.is_empty() {
y.yield_ok(Bytes::from(converted)).await;
}
}
let converted = normalizer.convert(&[], true).map_err(|source| Error::Generic {
store: "EcObjectStore",
source: Box::new(source),
})?;
if !converted.is_empty() {
y.yield_ok(Bytes::from(converted)).await;
}
Ok(())
})
.boxed()
}
#[cfg(test)]
mod stream_tests {
use super::*;
use std::sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
};
struct DropProbe(Arc<AtomicBool>);
impl Drop for DropProbe {
fn drop(&mut self) {
self.0.store(true, Ordering::SeqCst);
}
}
#[tokio::test]
async fn unicode_csv_quotes_drop_the_source_without_reading_ahead() {
let polls = Arc::new(AtomicUsize::new(0));
let dropped = Arc::new(AtomicBool::new(false));
let source =
futures::stream::unfold((DropProbe(Arc::clone(&dropped)), Arc::clone(&polls)), |(guard, polls)| async move {
polls.fetch_add(1, Ordering::SeqCst);
Some((Ok(Bytes::from_static("عvalueع\n".as_bytes())), (guard, polls)))
});
let csv = CsvSyntax {
quote: Some("ع"),
..Default::default()
};
let mut stream = normalize_csv_stream(source, &csv);
assert!(!stream.next().await.expect("first output").expect("valid CSV").is_empty());
assert_eq!(polls.load(Ordering::SeqCst), 1);
drop(stream);
assert!(dropped.load(Ordering::SeqCst), "cancellation must release the source reader");
assert_eq!(polls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn unicode_csv_quotes_preserve_source_errors_after_partial_output() {
let source = futures::stream::iter([
Ok(Bytes::from_static("عvalueع\n".as_bytes())),
Err(Error::Generic {
store: "fixture",
source: std::io::Error::other("source read failed").into(),
}),
]);
let csv = CsvSyntax {
quote: Some("ع"),
..Default::default()
};
let mut stream = normalize_csv_stream(source, &csv);
assert!(!stream.next().await.expect("partial output").expect("valid prefix").is_empty());
let error = stream
.next()
.await
.expect("source failure must remain visible")
.expect_err("must not return a successful tail");
assert!(error.to_string().contains("source read failed"));
assert!(stream.next().await.is_none());
}
}
+2
View File
@@ -23,12 +23,14 @@ use datafusion::{
use std::{error::Error as StdError, fmt::Display};
use thiserror::Error;
mod csv_input;
mod input_stream;
mod metrics;
pub mod object_store;
pub mod query;
pub mod server;
mod storage_api;
pub use csv_input::csv_input_requires_normalization;
pub use metrics::{SelectInputMetrics, SelectInputMetricsSnapshot};
pub use storage_api::SelectObjectSnapshot;
+109 -7
View File
@@ -64,6 +64,7 @@ use tokio::{io::AsyncReadExt, sync::OnceCell};
use tokio_util::io::ReaderStream;
use transform_stream::AsyncTryStream;
use crate::csv_input::{CsvSyntax, csv_input_requires_normalization, normalize_csv_stream};
use crate::storage_api::object_store::HTTPRangeSpec;
mod json_document;
@@ -345,6 +346,29 @@ impl EcObjectStore {
(self.need_convert || (delimiter.len() == 2 && delimiter != NORMALIZED_RECORD_DELIMITER)).then_some(delimiter)
}
fn convert_csv_stream<S>(&self, stream: S) -> BoxStream<'static, Result<Bytes>>
where
S: Stream<Item = Result<Bytes>> + Send + 'static,
{
if let Some(csv) = self.input.request.input_serialization.csv.as_ref()
&& csv_input_requires_normalization(csv.quote_character.as_deref(), csv.quote_escape_character.as_deref())
{
let syntax = CsvSyntax {
quote: csv.quote_character.as_deref(),
escape: csv.quote_escape_character.as_deref(),
field: csv.field_delimiter.as_deref(),
record: csv.record_delimiter.as_deref(),
comment: csv.comments.as_ref().and_then(|comment| comment.as_bytes().first().copied()),
};
return normalize_csv_stream(stream, &syntax);
}
convert_csv_delimiter_stream(
stream,
self.record_delimiter_for_conversion(),
self.need_convert.then(|| self.delimiter.clone()),
)
}
fn csv_has_header(&self) -> bool {
self.input
.request
@@ -820,7 +844,6 @@ impl ObjectStore for EcObjectStore {
});
}
let record_delimiter = self.record_delimiter_for_conversion();
let needs_scan_context = options.range.is_none() && has_effective_request_range;
let scan_context = if needs_scan_context {
if let Some(scan_range) = self.scan_range(original_size)? {
@@ -883,8 +906,7 @@ impl ObjectStore for EcObjectStore {
max_processed_bytes,
query_guard,
)?;
let stream =
convert_csv_delimiter_stream(stream, record_delimiter, self.need_convert.then(|| self.delimiter.clone()));
let stream = self.convert_csv_stream(stream);
GetResultPayload::Stream(stream)
}
} else if options.range.is_some() {
@@ -937,8 +959,7 @@ impl ObjectStore for EcObjectStore {
} else {
stream
};
let stream =
convert_csv_delimiter_stream(stream, record_delimiter, self.need_convert.then(|| self.delimiter.clone()));
let stream = self.convert_csv_stream(stream);
GetResultPayload::Stream(stream)
} else {
let stream_size = usize::try_from(original_size).map_err(|err| o_Error::Generic {
@@ -948,8 +969,7 @@ impl ObjectStore for EcObjectStore {
let stream = bytes_stream(ReaderStream::with_capacity(reader.stream, SELECT_DEFAULT_READ_BUFFER_SIZE), stream_size);
if meter_input {
let stream = meter_uncompressed_input_stream(stream, Arc::clone(&self.input_metrics));
let stream =
convert_csv_delimiter_stream(stream, record_delimiter, self.need_convert.then(|| self.delimiter.clone()));
let stream = self.convert_csv_stream(stream);
GetResultPayload::Stream(stream)
} else {
GetResultPayload::Stream(stream.boxed())
@@ -2866,6 +2886,88 @@ mod test {
assert_eq!(input_metrics.snapshot().bytes_processed, 2);
}
#[tokio::test]
async fn unicode_csv_quotes_preserve_raw_offsets_and_metrics() {
const BUCKET: &str = "s3select-unicode-csv-stream";
const HEADER: &str = "عnameع,عkindع\n";
const SKIP: &str = "عskipع,عzeroع\n";
const ROW: &str = "عA,Bع,عAععBع\n";
let data = format!("{HEADER}{SKIP}{ROW}");
let env = crate::storage_api::select_test_ecstore_env().await;
env.make_bucket(BUCKET, false).await;
for (object, compression, range_offset) in [
("plain.csv", None, None),
("range.csv", None, Some(0)),
("range-mid-character.csv", None, Some(1)),
("gzip.csv", Some(CompressionFormat::Gzip), None),
("bzip.csv", Some(CompressionFormat::Bzip2), None),
] {
let bytes = match compression {
Some(format) => encode_compressed_fixture(format, data.as_bytes()).await,
None => data.as_bytes().to_vec(),
};
let raw_size = bytes.len();
let mut reader = SelectPutObjReader::from_vec(bytes);
env.ecstore
.put_object(BUCKET, object, &mut reader, &Default::default())
.await
.expect("write Unicode CSV fixture");
let mut input = (*csv_input(BUCKET, object)).clone();
let csv = input.request.input_serialization.csv.as_mut().expect("CSV input");
csv.file_header_info = Some(FileHeaderInfo::from_static(FileHeaderInfo::USE));
csv.quote_character = Some("ع".to_owned());
csv.quote_escape_character = Some("\\".to_owned());
csv.record_delimiter = Some("\n".to_owned());
input.request.input_serialization.compression_type = compression.map(|format| {
CompressionType::from_static(match format {
CompressionFormat::Gzip => CompressionType::GZIP,
CompressionFormat::Bzip2 => CompressionType::BZIP2,
})
});
let start = HEADER.len() + SKIP.len();
if let Some(range_offset) = range_offset {
let offset = i64::try_from(start + range_offset).expect("fixture offset");
input.request.scan_range = Some(ScanRange {
start: Some(offset),
end: Some(offset),
});
}
let metrics = Arc::new(SelectInputMetrics::default());
let store = EcObjectStore::build_with_snapshot(
Arc::new(input),
Arc::new(GreedyMemoryPool::new(1024 * 1024)),
None,
Arc::clone(&metrics),
prepare_test_snapshot(BUCKET, object).await,
JsonSource::default(),
)
.expect("snapshot store");
let result = store
.get_opts(&Path::from(object), GetOptions::default())
.await
.expect("open Unicode CSV stream");
let GetResultPayload::Stream(stream) = result.payload else { panic!("CSV must remain streaming") };
let output = stream.try_collect::<Vec<_>>().await.expect("normalize CSV stream").concat();
let expected = match range_offset {
Some(0) => "\"name\",\"kind\"\n\"A,B\",\"AعB\"\n",
Some(_) => "\"name\",\"kind\"\n",
None => "\"name\",\"kind\"\n\"skip\",\"zero\"\n\"A,B\",\"AعB\"\n",
};
assert_eq!(output, expected.as_bytes(), "object={object}");
let measured = metrics.snapshot();
if let Some(range_offset) = range_offset {
// The range reader includes one byte of delimiter context and a
// separate header read; offsets always refer to the original CSV.
let processed = u64::try_from(ROW.len() + 1 - range_offset + HEADER.len()).expect("raw range length");
assert_eq!(measured.bytes_scanned, processed);
assert_eq!(measured.bytes_processed, processed);
} else {
assert_eq!(measured.bytes_scanned, u64::try_from(raw_size).expect("raw length"));
assert_eq!(measured.bytes_processed, u64::try_from(data.len()).expect("decoded length"));
}
}
}
#[tokio::test]
async fn compressed_object_uses_one_full_stream_and_rejects_internal_ranges() {
const BUCKET: &str = "s3select-compressed-object";
+25 -1
View File
@@ -456,7 +456,12 @@ impl SessionCtxFactory {
.is_some_and(|compression| compression.as_str() != CompressionType::NONE);
let metered_input_requires_single_file_scan =
input_metrics.is_some() && context.input.request.input_serialization.parquet.is_none();
let config = if custom_two_byte_record_delimiter
let normalized_csv_requires_single_file_scan =
context.input.request.input_serialization.csv.as_ref().is_some_and(|csv| {
crate::csv_input_requires_normalization(csv.quote_character.as_deref(), csv.quote_escape_character.as_deref())
});
let config = if normalized_csv_requires_single_file_scan
|| custom_two_byte_record_delimiter
|| scan_range_requires_single_file_scan
|| json_document_requires_single_file_scan
|| compressed_input_requires_single_file_scan
@@ -906,6 +911,25 @@ mod tests {
assert!(!session.inner().config().options().optimizer.repartition_file_scans);
}
#[tokio::test]
async fn unicode_csv_quotes_disable_file_scan_repartition() {
let mut context = test_context();
Arc::get_mut(&mut context.input)
.expect("unique context")
.request
.input_serialization
.csv
.as_mut()
.expect("CSV input")
.quote_character = Some("ع".to_owned());
let session = SessionCtxFactory::new(true)
.with_target_partitions(4)
.create_session_ctx(&context)
.await
.expect("Unicode CSV session");
assert!(!session.inner().config().options().optimizer.repartition_file_scans);
}
#[tokio::test]
async fn two_byte_csv_record_delimiter_disables_file_scan_repartition() {
let mut context = test_context();
+139 -23
View File
@@ -53,7 +53,7 @@ use rustfs_s3select_api::{
},
},
};
use s3s::dto::{CompressionType, FileHeaderInfo, JSONType, SelectObjectContentInput};
use s3s::dto::{FileHeaderInfo, JSONType, SelectObjectContentInput};
use std::sync::LazyLock;
use tokio::{
sync::Semaphore,
@@ -430,13 +430,6 @@ impl SimpleQueryDispatcher {
let path = format!("s3://{}/{}", self.input.bucket, self.input.key);
let table_path = ListingTableUrl::parse(path)?;
let compressed_input = self
.input
.request
.input_serialization
.compression_type
.as_ref()
.is_some_and(|compression| compression.as_str() != CompressionType::NONE);
let (listing_options, need_rename_volume_name, need_ignore_volume_name) =
if let Some(csv) = self.input.request.input_serialization.csv.as_ref() {
let mut need_rename_volume_name = false;
@@ -485,28 +478,27 @@ impl SimpleQueryDispatcher {
if let Some(quote) = csv.quote_character.as_ref() {
file_format = file_format.with_quote(quote.as_bytes().first().copied().unwrap_or_default());
}
if rustfs_s3select_api::csv_input_requires_normalization(
csv.quote_character.as_deref(),
csv.quote_escape_character.as_deref(),
) {
file_format = file_format
.with_quote(b'"')
.with_escape(None)
.with_delimiter(b',')
.with_terminator(Some(b'\n'))
.with_comment(None)
.with_newlines_in_values(true);
}
(
ListingOptions::new(Arc::new(file_format)).with_file_extension(if compressed_input {
EXACT_OBJECT_FILE_EXTENSION
} else {
".csv"
}),
ListingOptions::new(Arc::new(file_format)).with_file_extension(EXACT_OBJECT_FILE_EXTENSION),
need_rename_volume_name,
need_ignore_volume_name,
)
} else if self.input.request.input_serialization.json.is_some() {
let file_format = JsonFormat::default();
let file_extension = if compressed_input {
EXACT_OBJECT_FILE_EXTENSION.to_string()
} else {
std::path::Path::new(&self.input.key)
.extension()
.and_then(|extension| extension.to_str())
.map(|extension| format!(".{extension}"))
.unwrap_or_else(|| ".json".to_string())
};
(
ListingOptions::new(Arc::new(file_format)).with_file_extension(file_extension),
ListingOptions::new(Arc::new(file_format)).with_file_extension(EXACT_OBJECT_FILE_EXTENSION),
false,
false,
)
@@ -1531,6 +1523,130 @@ mod tests {
assert_eq!(error.select_error(), SelectError::InvalidDataSource);
}
#[tokio::test]
async fn unicode_csv_quotes_reach_arrow_without_changing_field_values() {
let cases = [
("ع", "\"", ",", "\n", "عcol1ع,عcol2ع,عcol3ع\n", vec![vec!["col1", "col2", "col3"]]),
(
"ع",
"\\",
",",
"\n",
"\"literal\",عA\\\"Bع,عAععBع\n",
vec![vec!["\"literal\"", "A\"B", "AعB"]],
),
("\"", "", ",", "\n", "\"A界\"B\",🦀\n", vec![vec!["A\"B", "🦀"]]),
("ع", "\\", "", "^Y", "عa界bع界عline\nbreakع^Y", vec![vec!["a界b", "line\nbreak"]]),
];
let env = snapshot_test_env().await;
for (index, (quote, escape, field, record, data, expected)) in cases.into_iter().enumerate() {
let mut input = test_input();
input.bucket = format!("select-unicode-quotes-{index}");
input.key = "records".to_owned();
let csv = input.request.input_serialization.csv.as_mut().expect("CSV input");
csv.file_header_info = Some(FileHeaderInfo::from_static(FileHeaderInfo::NONE));
csv.quote_character = Some(quote.to_owned());
csv.quote_escape_character = Some(escape.to_owned());
csv.field_delimiter = Some(field.to_owned());
csv.record_delimiter = Some(record.to_owned());
env.make_bucket(&input.bucket, false).await;
env.put_object_bytes(&input.bucket, &input.key, data.as_bytes().to_vec())
.await;
let snapshot = env.prepare_select_object_snapshot(&input.bucket, &input.key).await;
let input = Arc::new(input);
let dispatcher = production_dispatcher(Arc::clone(&input));
let query = Query::new_with_snapshot(QueryContext { input }, "SELECT * FROM S3Object".to_owned(), snapshot);
let output = dispatcher.execute_query(&query).await.expect("execute Unicode CSV query");
let mut stream = output.into_record_batch_stream().expect("record stream");
let mut rows = Vec::new();
while let Some(batch) = stream.next().await {
let batch = batch.expect("Arrow must receive valid UTF-8 fields");
for row in 0..batch.num_rows() {
rows.push(
batch
.columns()
.iter()
.map(|column| {
column
.as_any()
.downcast_ref::<StringArray>()
.expect("CSV string column")
.value(row)
.to_owned()
})
.collect::<Vec<_>>(),
);
}
}
assert_eq!(rows, expected, "fixture={index}");
}
}
#[tokio::test]
async fn select_uses_input_serialization_independently_of_object_extension() {
for (key, json) in [
("records", false),
("records.bin", false),
("records", true),
("records.csv", true),
] {
let mut input = test_input();
input.key = key.to_owned();
let data = if json {
input.request.input_serialization.csv = None;
input.request.input_serialization.json = Some(s3s::dto::JSONInput {
type_: Some(JSONType::from_static(JSONType::LINES)),
});
b"{\"value\":\"selected\"}\n".as_slice()
} else {
b"value\nselected\n".as_slice()
};
let input = Arc::new(input);
let optimizer = Arc::new(CascadeOptimizerBuilder::default().build());
let dispatcher = test_dispatcher_for_input(
Arc::clone(&input),
Arc::new(Semaphore::new(1)),
Duration::from_secs(30),
Arc::new(SqlQueryExecutionFactory::new(optimizer, Arc::new(LocalScheduler {}))),
);
let query = Query::new(QueryContext { input }, "SELECT * FROM S3Object".to_owned());
let machine = dispatcher.build_query_state_machine(query).await.expect("build query state");
let store_url = ObjectStoreUrl::parse("s3://test-bucket").expect("test store URL");
let store = machine
.session
.inner()
.runtime_env()
.object_store(&store_url)
.expect("test store");
store.put(&Path::from(key), data.into()).await.expect("write selected object");
store
.put(&Path::from(format!("{key}.other")), b"unrelated\nwrong\n".as_slice().into())
.await
.expect("write neighboring object");
let plan = dispatcher
.build_logical_plan(Arc::clone(&machine))
.await
.expect("infer schema without an extension filter")
.expect("select plan");
let output = dispatcher
.execute_logical_plan(plan, machine)
.await
.expect("execute selected object");
let mut stream = output.into_record_batch_stream().expect("record stream");
let mut values = Vec::new();
while let Some(batch) = stream.next().await {
let batch = batch.expect("selected batch");
let column = batch
.column(0)
.as_any()
.downcast_ref::<datafusion::arrow::array::StringArray>()
.expect("string column");
values.extend(column.iter().map(|value| value.expect("selected value").to_owned()));
}
assert_eq!(values, ["selected"], "key={key}, json={json}");
}
}
#[tokio::test]
async fn csv_query_uses_custom_record_delimiter_across_file_partitions() {
const ROW_COUNT: usize = 200_000;
@@ -67,6 +67,12 @@ Heal-side invariants that hold regardless of the caller:
- Read-repair's local TTL reservation dedups only its own source and does not block heals from other sources; the namespace lock is the backstop.
- The healing flag is never persisted, so there is no reverse risk of a leftover marker making a later commit yield incorrectly.
## Graceful root-heal restart recovery
Before a graceful shutdown cancels administrator cluster-wide heals, the manager saves unfinished requests on one coordinator disk as `.rustfs.sys/root-heal-<task-id>.json`. Startup replays the same task IDs and remaining execution budgets. Completion, cancellation, and replacement by `force_start` retire the record conditionally. An uncertain write or deletion does not create a fallback copy; an unsuccessful handoff retains the unclean-shutdown marker. Invalid or unsupported records remain on disk and defer root recovery without blocking the existing replacement-recovery path.
This handoff covers the same coordinator and storage topology while the record disk remains configured and readable. It does not migrate records when a pool is retired or provide failover after loss of that disk. Older versions do not understand these records; cancellation while downgraded cannot retire a newer version's pending record. If a terminal budget checkpoint cannot be written, the previous record is retained and a warning is logged; the remaining-budget guarantee requires that write to succeed. The format is separate from object metadata and erasure-set checkpoints.
## Regression tests
Both live in the test module of `crates/ecstore/src/set_disk/ops/heal.rs`:
@@ -38,6 +38,12 @@ Counts ignore blank lines and comments; compute them from the files. The lifecyc
"Supported" for the SSE row means RustFS encrypts and decrypts its own objects. MinIO SSE objects (SSE-S3, SSE-KMS, SSE-C) are not readable in default builds; see [minio-file-format-compat.md Part C](minio-file-format-compat.md#part-c--server-side-encryption-sse) for the `rio-v2` migration build.
### Client metadata expectations
`CopyObject` with `MetadataDirective=REPLACE` clears standard metadata fields that the request omits, including `Content-Type`; it does not retain the source type or infer a default. Clients requiring a MIME type on the copied object must send `Content-Type` with the replacement metadata. This contract is covered by `crates/e2e_test/src/copy_object_metadata_test.rs`. A client test that expects an implicit `application/octet-stream` does not match this behavior.
The MinIO-style `metadata=true` listing extension returns user metadata names without the HTTP `x-amz-meta-` prefix. It is not the standard S3 `ListObjectsV2` response. Clients that expect canonical HTTP header names in `UserMetadata` must normalize the names at that boundary; ordinary HEAD/GET metadata is unaffected. See `rustfs/src/app/bucket_usecase.rs` and its serialization tests.
## Replication Support Boundary
Site replication and bucket replication are not the same compatibility claim.
+34 -2
View File
@@ -689,8 +689,8 @@ fn normalize_input_serialization(input: &mut InputSerialization) -> S3Result<()>
));
}
validate_single_byte(csv.comments.as_deref(), S3ErrorCode::InvalidRequestParameter)?;
validate_single_byte(csv.quote_character.as_deref(), S3ErrorCode::InvalidRequestParameter)?;
validate_single_byte(csv.quote_escape_character.as_deref(), S3ErrorCode::InvalidRequestParameter)?;
validate_single_character(csv.quote_character.as_deref())?;
validate_single_character(csv.quote_escape_character.as_deref())?;
validate_input_record_delimiter(csv.record_delimiter.as_deref())?;
validate_input_delimiter_pair(csv.field_delimiter.as_deref(), csv.record_delimiter.as_deref())?;
}
@@ -778,6 +778,15 @@ fn invalid_scan_range_error() -> S3Error {
S3Error::with_message(S3ErrorCode::InvalidRequestParameter, INVALID_SCAN_RANGE_MESSAGE.to_string())
}
fn validate_single_character(value: Option<&str>) -> S3Result<()> {
if let Some(value) = value
&& value.chars().count() != 1
{
return Err(S3Error::new(S3ErrorCode::InvalidRequestParameter));
}
Ok(())
}
fn validate_single_byte(value: Option<&str>, code: S3ErrorCode) -> S3Result<()> {
if let Some(value) = value
&& value.len() != 1
@@ -3524,6 +3533,29 @@ mod tests {
assert_eq!(error.message(), Some(INVALID_SCAN_RANGE_MESSAGE));
}
#[test]
fn validate_accepts_single_unicode_csv_input_quotes() {
for quote in ["ع", "", "🦀"] {
let mut input = base_input();
let csv = input.request.input_serialization.csv.as_mut().expect("CSV input");
csv.quote_character = Some(quote.to_owned());
csv.quote_escape_character = Some(quote.to_owned());
validate_select_request(&HeaderMap::new(), &mut input).expect("one Unicode scalar is a valid CSV quote");
}
for quote in ["", "عع", "e\u{301}"] {
let mut input = base_input();
input
.request
.input_serialization
.csv
.as_mut()
.expect("CSV input")
.quote_character = Some(quote.to_owned());
let error = validate_select_request(&HeaderMap::new(), &mut input).expect_err("quote must be one scalar");
assert_eq!(error.code(), &S3ErrorCode::InvalidRequestParameter);
}
}
#[test]
fn validate_rejects_unknown_csv_header_mode_before_streaming() {
let mut input = base_input();
+17 -2
View File
@@ -280,6 +280,7 @@ pub(crate) async fn run_startup_shutdown_sequence(
let enable_scanner = get_env_bool_with_aliases(ENV_SCANNER_ENABLED, &[ENV_SCANNER_ENABLED_DEPRECATED], true);
let enable_heal = get_env_bool_with_aliases(ENV_HEAL_ENABLED, &[ENV_HEAL_ENABLED_DEPRECATED], true);
let mut heal_handoff_complete = true;
let background_steps = background_shutdown_steps(enable_scanner, enable_heal);
for step in &background_steps {
match step {
@@ -305,7 +306,19 @@ pub(crate) async fn run_startup_shutdown_sequence(
state = "stopping",
"Background service shutdown started"
);
shutdown_ahm_services();
if let Err(error) = shutdown_ahm_services().await {
heal_handoff_complete = false;
warn!(
target: "rustfs::main::handle_shutdown",
event = EVENT_BACKGROUND_SERVICE_SHUTDOWN,
component = LOG_COMPONENT_MAIN,
subsystem = LOG_SUBSYSTEM_STARTUP,
service = "ahm",
state = "handoff_failed",
error = %error,
"Heal shutdown handoff failed; retaining unclean-shutdown markers"
);
}
}
}
}
@@ -411,7 +424,9 @@ pub(crate) async fn run_startup_shutdown_sequence(
shutdown_optional_runtime_services(optional_runtime_shutdowns).await;
// The data plane is drained: record this shutdown as clean so the next
// startup skips the unclean-restart erasure-set heal.
rustfs_heal::heal::clear_unclean_shutdown_markers().await;
if heal_handoff_complete {
rustfs_heal::heal::clear_unclean_shutdown_markers().await;
}
state_manager.update(ServiceState::Stopped);
info!(
target: "rustfs::main::handle_shutdown",
+12 -2
View File
@@ -1006,12 +1006,22 @@ pub(crate) async fn apply_cors_headers(bucket: &str, method: &http::Method, head
}
// Access-Control-Allow-Headers (required for preflight if headers were requested)
if is_preflight && let Some(ref allowed_headers) = rule.allowed_headers {
let headers_str = allowed_headers.iter().map(|h| h.as_str()).collect::<Vec<_>>().join(", ");
if is_preflight && let Some(ref requested_headers) = requested_headers {
// Every requested header matched this rule; do not expose its wildcard
// or grant headers that the preflight did not request.
let headers_str = requested_headers.join(",");
if let Ok(headers_value) = HeaderValue::from_str(&headers_str) {
response_headers.insert(cors::response::ACCESS_CONTROL_ALLOW_HEADERS, headers_value);
}
}
if is_preflight {
let vary = if origin_reflected {
"Origin, Access-Control-Request-Method, Access-Control-Request-Headers"
} else {
"Access-Control-Request-Method, Access-Control-Request-Headers"
};
response_headers.insert(cors::standard::VARY, HeaderValue::from_static(vary));
}
// Access-Control-Expose-Headers (for actual requests)
if !is_preflight && let Some(ref expose_headers) = rule.expose_headers {
+4 -1
View File
@@ -1736,7 +1736,10 @@ mod tests {
"https://console.localhost",
);
assert_eq!(result.get(cors::response::ACCESS_CONTROL_ALLOW_CREDENTIALS).unwrap(), "true");
assert_eq!(result.get(cors::standard::VARY).unwrap(), "Origin");
assert_eq!(
result.get(cors::standard::VARY).unwrap(),
"Origin, Access-Control-Request-Method, Access-Control-Request-Headers"
);
set_bucket_metadata(bucket.to_string(), BucketMetadata::new(bucket))
.await
+1 -1
View File
@@ -26,7 +26,7 @@
6|crates/ecstore/src/config/com.rs
14|crates/ecstore/src/config/storageclass.rs
178|crates/ecstore/src/core/pools.rs
7|crates/ecstore/src/data_movement/mod.rs
6|crates/ecstore/src/data_movement/mod.rs
2|crates/ecstore/src/data_usage/local_snapshot.rs
12|crates/ecstore/src/data_usage/mod.rs
5|crates/ecstore/src/disk/local.rs