From 1650f3c2a650b4f0709327d3effd35571115730e Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sun, 6 Sep 2026 13:54:59 +0800 Subject: [PATCH] test(odm): verify global disable across restarts (#7268) * test(odm): cover global disable across restarts * test(odm): accept omitted V1 next marker * test(odm): gate backfill GET until the crash completes * test(odm): remove unused fault action import * test(odm): assert GET gate suspension without network races * test(odm): refresh verified Darwin E2E selection --- .config/e2e-full-selection.txt | 2 +- .config/e2e-smoke-selection.txt | 2 +- crates/e2e_test/src/fake_s3_target/mod.rs | 140 ++++++++++ .../on_demand_migration/interaction_test.rs | 242 +++++++++++++++++- 4 files changed, 381 insertions(+), 5 deletions(-) diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index 2cde94551..6a44dee52 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=a881fd7d3f5cb94654221ca85b8b30cce1b95e608824a55a15339cbc294e6d34 +sha256-darwin=53b05ac745905809d3828c6994bdd8ecf9d20b2b61a8a9d80fe15eb62f932193 sha256-linux=a2933d83dfe74ffa03410a0959333a1c48288b8469ca9f17273d449d7510c24b diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index 0a73fc814..6ebb3fcc0 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=a2542dc86bbff56b2177efc621785c56fa7e8d813b209b7d935e1e41a9f0ad15 +sha256=5db88c6fec94d4f269c7d9cfc128bd2adc27b3d7021127e2fa0b1daccc5f900f diff --git a/crates/e2e_test/src/fake_s3_target/mod.rs b/crates/e2e_test/src/fake_s3_target/mod.rs index 65f916fd9..24ca390ba 100644 --- a/crates/e2e_test/src/fake_s3_target/mod.rs +++ b/crates/e2e_test/src/fake_s3_target/mod.rs @@ -451,10 +451,40 @@ impl JournaledHeaders { struct ControlState { scripts: HashMap>, keyed_scripts: HashMap<(Operation, String), VecDeque>, + held_get: Option, requests: VecDeque, next_sequence: u64, } +#[derive(Clone)] +struct HeldGetObject { + bucket: String, + key: String, + entered: watch::Sender, + released: watch::Receiver, +} + +/// Holds every GET of one object, including retries, until this guard is dropped. +#[must_use = "dropping the guard releases the held GET requests"] +pub struct GetObjectGate { + control: Arc>, + entered: watch::Receiver, + released: watch::Sender, +} + +impl GetObjectGate { + pub async fn wait_until_entered(&mut self) -> Result<(), watch::error::RecvError> { + self.entered.wait_for(|count| *count > 0).await.map(|_| ()) + } +} + +impl Drop for GetObjectGate { + fn drop(&mut self) { + lock(&self.control).held_get = None; + self.released.send_replace(true); + } +} + #[derive(Default)] struct StoreState { assign_own_version_ids: bool, @@ -936,6 +966,30 @@ impl FakeS3Target { .extend(std::iter::repeat_n(action, times)); } + /// Hold one exact bucket/key before any GET response can reach the client. + /// The fixture supports one live gate; request and connection deadlines still apply. + pub fn hold_get_object(&self, bucket: &str, key: &str) -> GetObjectGate { + assert!( + bucket.len() <= MAX_RETAINED_IDENTIFIER_BYTES && key.len() <= MAX_RETAINED_IDENTIFIER_BYTES, + "held GET identifiers exceed the fixture limit" + ); + let mut state = lock(&self.control); + assert!(state.held_get.is_none(), "fake target already holds a GET gate"); + let (entered, entered_rx) = watch::channel(0); + let (released, released_rx) = watch::channel(false); + state.held_get = Some(HeldGetObject { + bucket: bucket.to_string(), + key: key.to_string(), + entered, + released: released_rx, + }); + GetObjectGate { + control: Arc::clone(&self.control), + entered: entered_rx, + released, + } + } + pub fn clear_faults(&self) { let mut state = lock(&self.control); state.scripts.clear(); @@ -2243,6 +2297,17 @@ impl S3 for FakeBackend { let fault = request_fault(&req); apply_non_body_fault(fault.as_ref(), &self.control).await?; let input = req.input; + let held_get = lock(&self.control) + .held_get + .as_ref() + .filter(|held| held.bucket == input.bucket && held.key == input.key) + .cloned(); + if let Some(mut held) = held_get { + held.entered.send_modify(|count| *count += 1); + // Keep the gate installed when a request is cancelled or times out: + // a retry must cross the same boundary before returning any bytes. + let _ = held.released.wait_for(|released| *released).await; + } let (version, versioned) = { let state = lock(&self.store); ( @@ -2894,6 +2959,81 @@ mod tests { aws_sdk_s3::primitives::DateTime::from_secs(4_102_444_800) } + #[tokio::test] + async fn get_object_gate_holds_retries_and_releases_on_drop() -> Result<(), BoxError> { + let target = FakeS3Target::start().await?; + let bucket = "gated-target"; + target.create_bucket(bucket); + for key in ["held", "unrelated"] { + target.put_seed_object(bucket, key, Bytes::from_static(b"payload"), &SeedMetadata::default()); + } + { + let gate = target.hold_get_object(bucket, "held"); + let request = || S3Request { + input: GetObjectInput { + bucket: bucket.to_string(), + key: "held".to_string(), + ..Default::default() + }, + method: Method::GET, + uri: Uri::from_static("/gated-target/held"), + headers: HeaderMap::new(), + extensions: http::Extensions::new(), + credentials: None, + region: None, + service: None, + trailing_headers: None, + }; + // Without a fault, only the gate can suspend this backend method. + let mut first = target.backend.get_object(request()); + assert!(futures::poll!(first.as_mut()).is_pending(), "the first GET must wait at the gate"); + drop(first); + let mut retry = target.backend.get_object(request()); + assert!(futures::poll!(retry.as_mut()).is_pending(), "a cancelled GET must not consume the gate"); + drop(gate); + let std::task::Poll::Ready(response) = futures::poll!(retry.as_mut()) else { + panic!("dropping the gate must release the waiting GET"); + }; + let mut body = response?.output.body.expect("released GET body"); + assert_eq!(body.next().await.transpose()?, Some(Bytes::from_static(b"payload"))); + assert!(body.next().await.is_none(), "released GET body must be complete"); + } + let client = client(&target); + let mut gate = target.hold_get_object(bucket, "held"); + let mut requests = tokio::task::JoinSet::new(); + let first = client.clone(); + requests.spawn(async move { get_bytes(&first, bucket, "held", None).await }); + timeout(Duration::from_secs(2), gate.wait_until_entered()).await??; + requests.abort_all(); + assert!( + requests + .join_next() + .await + .expect("first GET task") + .expect_err("cancel the first GET attempt") + .is_cancelled() + ); + + let retry = client.clone(); + requests.spawn(async move { get_bytes(&retry, bucket, "held", None).await }); + timeout(Duration::from_secs(2), gate.entered.wait_for(|count| *count == 2)).await??; + assert_eq!( + timeout(Duration::from_secs(2), get_bytes(&client, bucket, "unrelated", None)).await??, + Bytes::from_static(b"payload") + ); + assert!(requests.try_join_next().is_none(), "the retry must remain behind the gate"); + drop(gate); + assert_eq!( + timeout(Duration::from_secs(2), requests.join_next()) + .await? + .expect("retried GET task")??, + Bytes::from_static(b"payload") + ); + assert_eq!(get_bytes(&client, bucket, "held", None).await?, Bytes::from_static(b"payload")); + assert_eq!(target.count_requests(Operation::GetObject, "held"), 3); + Ok(()) + } + #[tokio::test] async fn object_lock_target_requires_a_checksum_on_locked_puts() -> Result<(), BoxError> { use aws_sdk_s3::error::ProvideErrorMetadata; diff --git a/crates/e2e_test/src/on_demand_migration/interaction_test.rs b/crates/e2e_test/src/on_demand_migration/interaction_test.rs index 93cafc05c..08a67adf8 100644 --- a/crates/e2e_test/src/on_demand_migration/interaction_test.rs +++ b/crates/e2e_test/src/on_demand_migration/interaction_test.rs @@ -22,8 +22,8 @@ //! local object and what the source was asked for. use super::common::{ - AdminResponse, BoxError, OdmEnvOptions, OdmSourceSpec, OdmTestEnv, SeedObject, start_configured_env, - start_configured_env_with, + ALLOW_LOOPBACK_SOURCE_ENV, AdminResponse, BackfillOp, BackfillRequest, BoxError, ODM_MODULE_SWITCH_ENV, ODM_SERVER_ENV, + OdmEnvOptions, OdmSourceSpec, OdmTestEnv, SeedObject, start_configured_env, start_configured_env_with, }; use crate::common::{RustFSTestEnvironment, replication_fast_env, signed_request}; use crate::fake_s3_target::{BucketMode, FAKE_ACCESS_KEY, FAKE_SECRET_KEY, FakeS3Target, Operation}; @@ -32,7 +32,7 @@ use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::types::{ BucketVersioningStatus, Event, FilterRule, FilterRuleName, NotificationConfiguration, NotificationConfigurationFilter, ObjectLockRetentionMode, QueueConfiguration, S3KeyFilter, ServerSideEncryption, ServerSideEncryptionByDefault, - ServerSideEncryptionConfiguration, ServerSideEncryptionRule, VersioningConfiguration, + ServerSideEncryptionConfiguration, ServerSideEncryptionRule, Tag, Tagging, VersioningConfiguration, }; use bytes::Bytes; use local_ip_address::local_ip; @@ -733,6 +733,242 @@ async fn test_odm_disable_keeps_pulled_objects_and_stops_source_traffic() -> Tes Ok(()) } +/// The process switch preserves configured buckets and unfinished jobs while +/// restoring local-only S3 behavior, including after an ordinary metadata write. +#[tokio::test] +async fn test_odm_global_disable_preserves_data_config_and_backfill_across_restarts() -> TestResult { + let bucket = "odm-global-disable"; + let mut env = start_configured_env(bucket, SOURCE_BUCKET, |spec| spec.policy.list_through = true).await?; + let pulled_key = "migrated/pulled.bin"; + let remote_key = "remote/untouched.bin"; + let pending_key = "backfill/pending.bin"; + let local_key = "local/kept.bin"; + let source_body = Bytes::from_static(b"source payload"); + let local_body = Bytes::from_static(b"client payload"); + env.seed_source( + SOURCE_BUCKET, + &[ + SeedObject::new(pulled_key, source_body.clone()), + SeedObject::new(remote_key, source_body.clone()), + SeedObject::new(pending_key, source_body.clone()), + ], + ); + env.client + .put_object() + .bucket(bucket) + .key(local_key) + .body(local_body.clone().into()) + .send() + .await?; + let pulled = env.raw_get(bucket, pulled_key).await?; + assert_eq!(pulled.status, 200); + assert_eq!(pulled.header(ODM_RESPONSE_HEADER), Some("source")); + assert_eq!(pulled.body, source_body); + let stored = env.raw_get(bucket, pulled_key).await?; + assert_eq!(stored.status, 200); + assert_eq!(stored.header(ODM_RESPONSE_HEADER), None, "the inline pull has committed locally"); + assert_eq!(stored.body, source_body); + let config = env.get_config(bucket).await?; + assert_eq!(config.status, 200, "{}", config.body); + let config = config.json()?; + + // Hold every attempt until the process has exited, so retries cannot commit + // the only backfill object before the crash. The start checkpoint exists. + let mut pending_get = env.source.hold_get_object(SOURCE_BUCKET, pending_key); + let started = env + .start_backfill( + bucket, + BackfillRequest { + prefix: Some("backfill/".to_string()), + ..BackfillRequest::default() + }, + ) + .await?; + assert_eq!(started.status, 200, "{}", started.body); + let job_id = started.json()?["job"]["job_id"].as_str().ok_or("missing job ID")?.to_string(); + tokio::time::timeout(Duration::from_secs(10), pending_get.wait_until_entered()) + .await + .expect("backfill never reached the held source GET")?; + let process = env.rustfs.process.as_mut().ok_or("missing RustFS process before crash")?; + assert!(process.try_wait()?.is_none(), "RustFS exited before the controlled crash"); + process.kill()?; + let stopped = process.wait()?; + assert!(!stopped.success(), "the interrupted process must exit after being killed"); + drop(env.rustfs.process.take()); + drop(pending_get); + env.source.take_requests(); + env.rustfs + .restart_server_preserving_data(vec![], &[(ODM_MODULE_SWITCH_ENV, "false"), (ALLOW_LOOPBACK_SOURCE_ENV, "true")]) + .await?; + + let off_config = env.get_config(bucket).await?; + assert_eq!(off_config.status, 200, "{}", off_config.body); + assert_eq!(off_config.json()?, config, "the saved configuration and timestamp survive disabling"); + let status = env.status_json(bucket).await?; + assert_eq!(status["configured"], true, "{status}"); + assert_eq!(status["enabled"], true, "the bucket remains configured as enabled: {status}"); + assert_eq!(status["module_enabled"], false, "{status}"); + assert_eq!(status["counters"], Value::Null, "no bucket runtime is installed: {status}"); + let checkpoint = env.backfill_job(bucket).await?.ok_or("disabled module lost the checkpoint")?; + assert_eq!(checkpoint["job_id"], job_id); + assert_eq!(checkpoint["state"], "running", "the interrupted job is retained: {checkpoint}"); + + for (key, body) in [(local_key, &local_body), (pulled_key, &source_body)] { + let get = env.raw_get(bucket, key).await?; + assert_eq!(get.status, 200); + assert_eq!(&get.body, body); + assert_eq!(get.header(ODM_RESPONSE_HEADER), None); + let head = env.client.head_object().bucket(bucket).key(key).send().await?; + assert_eq!(head.content_length(), Some(i64::try_from(body.len())?)); + } + for key in [remote_key, pending_key] { + let get = env.raw_get(bucket, key).await?; + assert_eq!(get.status, 404, "disabled source GET {key}: {}", String::from_utf8_lossy(&get.body)); + let head = env.client.head_object().bucket(bucket).key(key).send().await; + let err = head.expect_err("a source-only object must remain absent locally"); + assert_eq!(err.raw_response().map(|response| response.status().as_u16()), Some(404)); + } + + let replacement = Bytes::from_static(b"written while the module is off"); + for key in [local_key, "local/deleted.bin"] { + env.client + .put_object() + .bucket(bucket) + .key(key) + .body(replacement.clone().into()) + .send() + .await?; + } + env.client + .delete_object() + .bucket(bucket) + .key("local/deleted.bin") + .send() + .await?; + assert_eq!(env.raw_get(bucket, "local/deleted.bin").await?.status, 404); + assert_eq!(env.raw_get(bucket, local_key).await?.body, replacement); + + // Both wire protocols must finish their local pages even though the saved + // configuration still requests list-through. + for use_v2 in [false, true] { + let mut cursor = None; + let mut listed = Vec::new(); + for page_number in 0..2 { + let (keys, truncated, next) = if use_v2 { + let page = env + .client + .list_objects_v2() + .bucket(bucket) + .max_keys(1) + .set_continuation_token(cursor) + .send() + .await?; + ( + page.contents() + .iter() + .map(|object| object.key().expect("listed key").to_string()) + .collect::>(), + page.is_truncated(), + page.next_continuation_token().map(str::to_string), + ) + } else { + let page = env + .client + .list_objects() + .bucket(bucket) + .max_keys(1) + .set_marker(cursor) + .send() + .await?; + // V1 may omit NextMarker without a delimiter; clients then + // continue from the last returned key. + let next = page.next_marker().or_else(|| { + if page.is_truncated() == Some(true) { + page.contents().last().and_then(|object| object.key()) + } else { + None + } + }); + ( + page.contents() + .iter() + .map(|object| object.key().expect("listed key").to_string()) + .collect::>(), + page.is_truncated(), + next.map(str::to_string), + ) + }; + assert_eq!(keys.len(), 1, "one local key per page, V2={use_v2}"); + assert_eq!(truncated, Some(page_number == 0), "local pagination must terminate, V2={use_v2}"); + if page_number == 0 { + assert!(next.as_ref().is_some_and(|value| !value.is_empty()), "missing local cursor, V2={use_v2}"); + } + cursor = next; + listed.extend(keys); + } + assert_eq!(listed, [local_key, pulled_key], "source-only keys must stay absent, V2={use_v2}"); + } + + let spec = env.fake_source_spec(SOURCE_BUCKET); + for response in [ + env.configure_source(bucket, &spec).await?, + env.validate_source(bucket, &spec).await?, + env.backfill(bucket, BackfillOp::Start(BackfillRequest::default())).await?, + ] { + assert_eq!(response.status, 400, "{}", response.body); + assert!(response.body.contains("OnDemandMigrationDisabled"), "{}", response.body); + } + let tagging = Tagging::builder() + .tag_set(Tag::builder().key("module").value("disabled").build()?) + .build()?; + env.client + .put_bucket_tagging() + .bucket(bucket) + .tagging(tagging.clone()) + .send() + .await?; + assert_eq!(env.get_config(bucket).await?.json()?, config, "an unrelated metadata write preserves ODM"); + assert_eq!( + env.backfill_job(bucket).await?, + Some(checkpoint), + "no recovery or checkpoint update while disabled" + ); + assert!( + env.source.requests().is_empty(), + "disabled startup and all requests must leave the source untouched" + ); + + env.rustfs.restart_server_preserving_data(vec![], ODM_SERVER_ENV).await?; + env.wait_until_source_consulted(bucket).await?; + assert_eq!( + env.get_config(bucket).await?.json()?, + config, + "reenabling uses the persisted configuration" + ); + let tags = env.client.get_bucket_tagging().bucket(bucket).send().await?; + assert_eq!(tags.tag_set(), tagging.tag_set(), "the ordinary metadata write also persists"); + let resumed = env.raw_get(bucket, remote_key).await?; + assert_eq!(resumed.status, 200); + assert_eq!(resumed.header(ODM_RESPONSE_HEADER), Some("source")); + assert_eq!(resumed.body, source_body, "stored credentials still authenticate without reconfiguration"); + let completed = env + .wait_for_backfill(bucket, SETTLE, |job| job["state"] == "completed") + .await?; + assert_eq!(completed["job_id"], job_id, "the interrupted job resumes without a new start"); + assert_eq!(completed["failed"], 0, "{completed}"); + for (key, body) in [ + (local_key, &replacement), + (pulled_key, &source_body), + (pending_key, &source_body), + ] { + let get = env.raw_get(bucket, key).await?; + assert_eq!(get.status, 200); + assert_eq!(&get.body, body); + assert_eq!(get.header(ODM_RESPONSE_HEADER), None, "{key} remains stored locally"); + } + Ok(()) +} + /// Case 19: the admin surface an operator sees — the configuration read back /// without its secret, and a status document whose counters match the source /// journal exactly.