From e24eae9eaad916ecf8a1e546179674409a9b5027 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 9 Sep 2026 08:33:59 +0800 Subject: [PATCH] fix: address release acceptance regressions (#7541) --- .config/e2e-full-selection.txt | 4 +- .config/e2e-nightly-selection.txt | 4 +- .../e2e_test/src/cluster_concurrency_test.rs | 55 +++ crates/e2e_test/src/distributed/harness.rs | 21 - .../s3_during_data_movement_test.rs | 30 +- .../src/heal_erasure_disk_rebuild_test.rs | 27 +- .../src/kms/kms_fault_recovery_test.rs | 4 +- crates/ecstore/src/core/pools.rs | 448 ++++++++++++++++- crates/ecstore/src/data_movement/mod.rs | 33 +- crates/ecstore/src/store/list_objects.rs | 257 +++++++++- crates/heal/src/heal/manager.rs | 147 +++++- crates/heal/src/heal/manager/root_recovery.rs | 342 +++++++++++++ crates/heal/src/heal/manager/scheduler.rs | 40 ++ crates/heal/src/heal/manager/tests.rs | 2 + .../src/heal/manager/tests/root_recovery.rs | 465 ++++++++++++++++++ crates/heal/src/heal/task.rs | 5 + crates/heal/src/lib.rs | 6 +- crates/s3select-api/src/csv_input.rs | 432 ++++++++++++++++ crates/s3select-api/src/lib.rs | 2 + crates/s3select-api/src/object_store.rs | 116 ++++- crates/s3select-api/src/query/session.rs | 26 +- .../s3select-query/src/dispatcher/manager.rs | 162 +++++- docs/architecture/heal-concurrency-model.md | 6 + docs/architecture/s3-compatibility-matrix.md | 6 + rustfs/src/app/select_object.rs | 36 +- rustfs/src/startup_shutdown.rs | 19 +- rustfs/src/storage/ecfs_extend.rs | 14 +- rustfs/src/storage/ecfs_test.rs | 5 +- scripts/error-other-format-baseline.txt | 2 +- 29 files changed, 2582 insertions(+), 134 deletions(-) create mode 100644 crates/heal/src/heal/manager/root_recovery.rs create mode 100644 crates/heal/src/heal/manager/tests/root_recovery.rs create mode 100644 crates/s3select-api/src/csv_input.rs diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index 349e17819..ceb69305a 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=f0c78fdb93471575d9a64c5c46eae6c806bdd0bc10a6e33d7fb574aabd8db5a3 -sha256-linux=03ed7016cab672de9320e31375a0358eceacb4408b0e79cf063614fa7c878b87 +sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5 +sha256-linux=9515861be899ceb10e2e0ef93c34208bb7a7a8a7f8067a02db4cfba23270ebd6 diff --git a/.config/e2e-nightly-selection.txt b/.config/e2e-nightly-selection.txt index 1bb5125c9..bec86f799 100644 --- a/.config/e2e-nightly-selection.txt +++ b/.config/e2e-nightly-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=364f2329a7b72eb9f1608dbe1a3af37af4095354014f3cbe23ca448492d89961 -sha256-linux=60983f1ebe7068cf660d473c5f76c76a650410ccc99d71934ddca7fd67607987 +sha256-darwin=83a7dcaffd5a789517ae9f02a224f66a9713937885cff96fca2ad7e216f197ae +sha256-linux=626c10f8c964507ff987b6c86069e9019dc6d2ae7fb02db9be5df5aa8cc5145b diff --git a/crates/e2e_test/src/cluster_concurrency_test.rs b/crates/e2e_test/src/cluster_concurrency_test.rs index 03d73954a..177cdbd8f 100644 --- a/crates/e2e_test/src/cluster_concurrency_test.rs +++ b/crates/e2e_test/src/cluster_concurrency_test.rs @@ -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) diff --git a/crates/e2e_test/src/distributed/harness.rs b/crates/e2e_test/src/distributed/harness.rs index aed1021dc..bf53789c9 100644 --- a/crates/e2e_test/src/distributed/harness.rs +++ b/crates/e2e_test/src/distributed/harness.rs @@ -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, diff --git a/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs b/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs index 66c835bbd..79b7710d1 100644 --- a/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs +++ b/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs @@ -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(()) } diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index f8763d482..c9df49d4a 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -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(); diff --git a/crates/e2e_test/src/kms/kms_fault_recovery_test.rs b/crates/e2e_test/src/kms/kms_fault_recovery_test.rs index 2e66c9bf8..4fed36609 100644 --- a/crates/e2e_test/src/kms/kms_fault_recovery_test.rs +++ b/crates/e2e_test/src/kms/kms_fault_recovery_test.rs @@ -109,10 +109,10 @@ async fn test_kms_key_directory_unavailable() -> Result<(), Box 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(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)); diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index b89f806bd..f4a590bf2 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -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] diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 82e6b1aa2..459ef0870 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -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"); diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 732778d15..23b586280 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -841,6 +841,10 @@ pub struct HealManager { replacement_recovery_anchors: Arc>>, /// Set IDs whose durable replacement metadata is corrupt or conflicting. replacement_recovery_blocked_sets: Arc>>, + /// Durable handoff of interrupted administrator root traversals. + root_recovery: Arc, + /// Keep forceStart's cancellation side effects inside the shutdown fence. + force_start_shutdown: Mutex<()>, /// Storage layer interface storage: Arc, /// Cancel token @@ -876,6 +880,7 @@ struct HealQueueContext<'a> { retrying_heals: &'a Arc>>, mrf_repair_notice_targets: &'a Arc>>>, replacement_recovery_anchors: &'a Arc>>, + root_recovery: &'a Arc, config: &'a Arc>, statistics: &'a Arc>, storage: &'a Arc, @@ -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 = { 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::>(); + 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::>(); 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; diff --git a/crates/heal/src/heal/manager/root_recovery.rs b/crates/heal/src/heal/manager/root_recovery.rs new file mode 100644 index 000000000..f32ff1028 --- /dev/null +++ b/crates/heal/src/heal/manager/root_recovery.rs @@ -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>, +} + +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 { + 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 { + 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 { + 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) -> Self { + Self { + mutation: Mutex::new(()), + disks: Some(disks), + } + } + + async fn disks(&self) -> Result> { + #[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::>(); + disks.sort_by_key(|disk| EcstoreDiskAPI::endpoint(disk.as_ref()).to_string()); + Ok(disks) + } + + async fn find(disks: &[DiskStore], task_id: &str) -> Result> { + 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 { + 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 { + 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> { + 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(()) + } +} diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index 1419f5587..a3820fa2b 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -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 { diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 6e551b829..40ff0e701 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -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, diff --git a/crates/heal/src/heal/manager/tests/root_recovery.rs b/crates/heal/src/heal/manager/tests/root_recovery.rs new file mode 100644 index 000000000..28fb83288 --- /dev/null +++ b/crates/heal/src/heal/manager/tests/root_recovery.rs @@ -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) -> 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 { + 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::>(); + 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::>(); + 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::>(); + 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"); + } + } +} diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 816f22afd..75996ed90 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -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 { let mut request = self.retry_request(); if self.options.timeout.is_some() { diff --git a/crates/heal/src/lib.rs b/crates/heal/src/lib.rs index 733701e45..38a3883d7 100644 --- a/crates/heal/src/lib.rs +++ b/crates/heal/src/lib.rs @@ -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 { diff --git a/crates/s3select-api/src/csv_input.rs b/crates/s3select-api/src/csv_input.rs new file mode 100644 index 000000000..b63d82d84 --- /dev/null +++ b/crates/s3select-api/src/csv_input.rs @@ -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, +} + +#[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, + escape: Vec, + field: Vec, + record: Vec, + comment: Option, + default_records: bool, + state: State, + record_start: bool, + carry: Vec, + 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, 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, 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 { + 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\0dع\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("e_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(stream: S, csv: &CsvSyntax<'_>) -> BoxStream<'static, Result> +where + S: Stream> + Send + 'static, +{ + let mut normalizer = CsvInputNormalizer::new(csv); + AsyncTryStream::::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); + + 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()); + } +} diff --git a/crates/s3select-api/src/lib.rs b/crates/s3select-api/src/lib.rs index a443a07a1..cebec28a6 100644 --- a/crates/s3select-api/src/lib.rs +++ b/crates/s3select-api/src/lib.rs @@ -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; diff --git a/crates/s3select-api/src/object_store.rs b/crates/s3select-api/src/object_store.rs index 499436333..c2b2240d1 100644 --- a/crates/s3select-api/src/object_store.rs +++ b/crates/s3select-api/src/object_store.rs @@ -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(&self, stream: S) -> BoxStream<'static, Result> + where + S: Stream> + 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::>().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"; diff --git a/crates/s3select-api/src/query/session.rs b/crates/s3select-api/src/query/session.rs index 7df3cf60a..124c0c85d 100644 --- a/crates/s3select-api/src/query/session.rs +++ b/crates/s3select-api/src/query/session.rs @@ -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(); diff --git a/crates/s3select-query/src/dispatcher/manager.rs b/crates/s3select-query/src/dispatcher/manager.rs index 05a00e3e5..22a733fe7 100644 --- a/crates/s3select-query/src/dispatcher/manager.rs +++ b/crates/s3select-query/src/dispatcher/manager.rs @@ -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::() + .expect("CSV string column") + .value(row) + .to_owned() + }) + .collect::>(), + ); + } + } + 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::() + .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; diff --git a/docs/architecture/heal-concurrency-model.md b/docs/architecture/heal-concurrency-model.md index f907b44fe..d7f2368d1 100644 --- a/docs/architecture/heal-concurrency-model.md +++ b/docs/architecture/heal-concurrency-model.md @@ -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-.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`: diff --git a/docs/architecture/s3-compatibility-matrix.md b/docs/architecture/s3-compatibility-matrix.md index f296003dd..05b57c074 100644 --- a/docs/architecture/s3-compatibility-matrix.md +++ b/docs/architecture/s3-compatibility-matrix.md @@ -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. diff --git a/rustfs/src/app/select_object.rs b/rustfs/src/app/select_object.rs index 48ea1ce50..9e67b90ec 100644 --- a/rustfs/src/app/select_object.rs +++ b/rustfs/src/app/select_object.rs @@ -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(); diff --git a/rustfs/src/startup_shutdown.rs b/rustfs/src/startup_shutdown.rs index 70b1a0701..50f855662 100644 --- a/rustfs/src/startup_shutdown.rs +++ b/rustfs/src/startup_shutdown.rs @@ -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", diff --git a/rustfs/src/storage/ecfs_extend.rs b/rustfs/src/storage/ecfs_extend.rs index 9ecd207a9..c3807527e 100644 --- a/rustfs/src/storage/ecfs_extend.rs +++ b/rustfs/src/storage/ecfs_extend.rs @@ -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::>().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 { diff --git a/rustfs/src/storage/ecfs_test.rs b/rustfs/src/storage/ecfs_test.rs index cf4f011a4..87852378c 100644 --- a/rustfs/src/storage/ecfs_test.rs +++ b/rustfs/src/storage/ecfs_test.rs @@ -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 diff --git a/scripts/error-other-format-baseline.txt b/scripts/error-other-format-baseline.txt index b21e757c4..7b7abe72a 100644 --- a/scripts/error-other-format-baseline.txt +++ b/scripts/error-other-format-baseline.txt @@ -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