diff --git a/crates/ecstore/src/ecstore_validation_blackbox.rs b/crates/ecstore/src/ecstore_validation_blackbox.rs index 3c94e3e9b..31442dde1 100644 --- a/crates/ecstore/src/ecstore_validation_blackbox.rs +++ b/crates/ecstore/src/ecstore_validation_blackbox.rs @@ -20,7 +20,7 @@ use tokio::sync::RwLock; /// Returns the backing [`tempfile::TempDir`]s alongside the set so callers keep /// them alive for the test's duration and the directories are removed on drop. -async fn make_local_set_disks(drive_count: usize, parity_count: usize) -> (Vec, Arc) { +pub(crate) async fn make_local_set_disks(drive_count: usize, parity_count: usize) -> (Vec, Arc) { let format = FormatV3::new(1, drive_count); let mut dirs = Vec::with_capacity(drive_count); let mut endpoints = Vec::with_capacity(drive_count); diff --git a/crates/ecstore/src/lib.rs b/crates/ecstore/src/lib.rs index 65340afba..e6c61e9a1 100644 --- a/crates/ecstore/src/lib.rs +++ b/crates/ecstore/src/lib.rs @@ -87,7 +87,7 @@ mod rio_tests { } #[cfg(test)] -mod ecstore_validation_blackbox; +pub(crate) mod ecstore_validation_blackbox; #[cfg(test)] pub(crate) mod test_metrics; diff --git a/crates/ecstore/src/object_api/body_cache_hook.rs b/crates/ecstore/src/object_api/body_cache_hook.rs index 14fd1e7e7..d7cb265f6 100644 --- a/crates/ecstore/src/object_api/body_cache_hook.rs +++ b/crates/ecstore/src/object_api/body_cache_hook.rs @@ -22,7 +22,7 @@ use crate::object_api::ObjectInfo; use bytes::Bytes; -use std::sync::{Arc, OnceLock}; +use std::sync::{Arc, RwLock}; /// Serves full-object GET bodies from a cache keyed by object identity. /// @@ -35,15 +35,48 @@ pub trait GetObjectBodyCacheHook: Send + Sync + 'static { async fn lookup(&self, bucket: &str, object: &str, info: &ObjectInfo) -> Option; } -static GET_OBJECT_BODY_CACHE_HOOK: OnceLock> = OnceLock::new(); +// `RwLock>>` rather than `ArcSwapOption`: arc-swap's +// `RefCnt` is implemented only for the sized `Arc` (it stores a thin +// `*mut T`), so it cannot hold an `Arc` without a +// sized newtype wrapper. The probe reads this slot once per full-object GET, +// but the read guard only clones an `Arc`, which is negligible next to the +// metadata quorum fan-out already completed before the probe. Registration is +// a startup / config-reload event, so writer contention is a non-issue. +static GET_OBJECT_BODY_CACHE_HOOK: RwLock>> = RwLock::new(None); -/// Register the process-wide GET body cache hook. First registration wins; -/// later calls are ignored so tests and re-inits cannot swap the hook midway. +/// Register (or re-register) the process-wide GET body cache hook. +/// +/// Re-registration atomically swaps to `hook`, so a rebuilt `AppContext` +/// (config reload, or the test re-init pattern) leaves ecstore's GET probe +/// pointed at the newest adapter. A first-wins slot would instead pin the probe +/// to the original adapter while every usecase-layer fill and invalidation +/// targeted the replacement, silently degrading the feature to a 0% hit rate +/// and stranding entries in the unreachable cache until their TTL (backlog#1126). +/// +/// Replacing a *different* hook instance is logged at WARN: in production the +/// hook is installed exactly once per process, so a swap to a distinct instance +/// signals an unexpected re-init and orphans the previous adapter's cache. pub fn register_get_object_body_cache_hook(hook: Arc) { - let _ = GET_OBJECT_BODY_CACHE_HOOK.set(hook); + let mut slot = GET_OBJECT_BODY_CACHE_HOOK.write().unwrap_or_else(|e| e.into_inner()); + if let Some(previous) = slot.as_ref() + && !Arc::ptr_eq(previous, &hook) + { + tracing::warn!( + "GET object body cache hook re-registered with a different instance; \ + the previous adapter's cache is now unreachable by ecstore's GET probe" + ); + } + *slot = Some(hook); } /// The registered hook, if any. -pub(crate) fn get_object_body_cache_hook() -> Option<&'static Arc> { - GET_OBJECT_BODY_CACHE_HOOK.get() +pub(crate) fn get_object_body_cache_hook() -> Option> { + GET_OBJECT_BODY_CACHE_HOOK.read().unwrap_or_else(|e| e.into_inner()).clone() +} + +/// Test-only: unregister the hook so tests can register and clear the slot +/// deterministically without leaking a hook into unrelated tests. +#[cfg(test)] +pub(crate) fn clear_get_object_body_cache_hook() { + *GET_OBJECT_BODY_CACHE_HOOK.write().unwrap_or_else(|e| e.into_inner()) = None; } diff --git a/crates/ecstore/src/object_api/mod.rs b/crates/ecstore/src/object_api/mod.rs index 67fb33229..d08cc4285 100644 --- a/crates/ecstore/src/object_api/mod.rs +++ b/crates/ecstore/src/object_api/mod.rs @@ -56,6 +56,8 @@ mod body_cache_hook; mod readers; mod types; +#[cfg(test)] +pub(crate) use body_cache_hook::clear_get_object_body_cache_hook; pub(crate) use body_cache_hook::get_object_body_cache_hook; pub use body_cache_hook::{GetObjectBodyCacheHook, register_get_object_body_cache_hook}; pub use readers::*; diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 7294509db..e99b8373a 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -3108,3 +3108,247 @@ mod body_cache_hook_gate_tests { assert_eq!(full_object_plaintext_len(&None, &ObjectOptions::default(), &info), None); } } + +#[cfg(test)] +mod body_cache_hook_e2e_tests { + //! End-to-end regression coverage for the two merged P0 data-corruption + //! fixes, driving the *real* `get_object_reader` gate + probe + reader + //! construction rather than the `full_object_plaintext_len` predicate in + //! isolation. The predicate-only tests above cannot catch a caller that + //! opens a new shortcut serving the cached body directly — which was the + //! original form of both P0s (backlog#1108 / #1109) and the reason ODC-21 + //! (backlog#1126) made the hook re-registrable so these tests can install a + //! deterministic hook. + //! + //! A stand-in `GetObjectBodyCacheHook` plays the app-layer cache: the + //! adapter itself lives above ecstore and cannot be reached from here, but + //! the hook trait is exactly the production injection point, so exercising + //! it through `get_object_reader` is a true end-to-end test of the ecstore + //! side (gate decision -> probe -> served bytes and published size). + //! + //! The tests share one process-global hook slot, so they serialize under a + //! single key and clear the slot on the way out; each uses a unique + //! bucket/object and the hook returns `None` for anything else, so a stray + //! concurrent GET in another test is unaffected. + + use crate::ecstore_validation_blackbox::make_local_set_disks; + use crate::io_support::rio::{HashReader, compression_metadata_value, compression_reader}; + use crate::object_api::{ + GetObjectBodyCacheHook, ObjectInfo, ObjectOptions, PutObjReader, clear_get_object_body_cache_hook, + register_get_object_body_cache_hook, + }; + use crate::set_disk::SetDisks; + use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; + use crate::storage_api_contracts::object::ObjectIO as _; + use bytes::Bytes; + use http::HeaderMap; + use rustfs_utils::CompressionAlgorithm; + use rustfs_utils::http::{SUFFIX_ACTUAL_SIZE, SUFFIX_COMPRESSION, insert_str}; + use std::collections::HashMap; + use std::io::Cursor; + use std::sync::Arc; + use tokio::io::AsyncReadExt as _; + + /// A stand-in for the app-layer object-data cache: returns `body` only for + /// the one primed key, mirroring how the real adapter answers a probe. + struct PrimedBodyHook { + bucket: String, + object: String, + body: Bytes, + } + + #[async_trait::async_trait] + impl GetObjectBodyCacheHook for PrimedBodyHook { + async fn lookup(&self, bucket: &str, object: &str, _info: &ObjectInfo) -> Option { + (bucket == self.bucket && object == self.object).then(|| self.body.clone()) + } + } + + /// Registers a primed hook and clears it on drop so no hook leaks into an + /// unrelated test sharing the process-global slot. + struct HookGuard; + impl HookGuard { + fn install(bucket: &str, object: &str, body: Bytes) -> Self { + register_get_object_body_cache_hook(Arc::new(PrimedBodyHook { + bucket: bucket.to_string(), + object: object.to_string(), + body, + })); + HookGuard + } + } + impl Drop for HookGuard { + fn drop(&mut self) { + clear_get_object_body_cache_hook(); + } + } + + /// Compresses `plaintext` with the same codec the real PUT path uses, so + /// the stored bytes round-trip through `get_object_reader`'s decompressor. + async fn compress(plaintext: &[u8]) -> Vec { + let mut reader = compression_reader(Cursor::new(plaintext.to_vec()), CompressionAlgorithm::default(), false); + let mut compressed = Vec::new(); + reader.read_to_end(&mut compressed).await.expect("compress plaintext"); + compressed + } + + /// Writes a genuinely compressed object: the stored data is `compressed` + /// and the metadata marks it compressed with `plaintext.len()` as the + /// decompressed length, exactly as the app-layer compress path records it. + async fn put_compressed_object(set_disks: &Arc, bucket: &str, object: &str, plaintext: &[u8], compressed: &[u8]) { + let mut user_defined = HashMap::new(); + insert_str( + &mut user_defined, + SUFFIX_COMPRESSION, + compression_metadata_value(CompressionAlgorithm::default()), + ); + insert_str(&mut user_defined, SUFFIX_ACTUAL_SIZE, plaintext.len().to_string()); + + let opts = ObjectOptions { + no_lock: true, + user_defined, + ..Default::default() + }; + let stream = HashReader::from_stream( + Cursor::new(compressed.to_vec()), + compressed.len() as i64, // stored (compressed) size + plaintext.len() as i64, // actual (decompressed) size + None, + None, + false, + ) + .expect("hash reader over compressed bytes"); + let mut reader = PutObjReader::new(stream); + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + set_disks + .put_object(bucket, object, &mut reader, &opts) + .await + .expect("compressed object should be written"); + } + + async fn read_to_vec(set_disks: &Arc, bucket: &str, object: &str, opts: &ObjectOptions) -> (Vec, i64) { + let mut reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), opts) + .await + .expect("object reader should open"); + let published_size = reader.object_info.size; + let mut body = Vec::new(); + reader.stream.read_to_end(&mut body).await.expect("object should stream"); + (body, published_size) + } + + /// A compressible payload large enough to guarantee stored != plaintext and + /// to keep the object off the inline / direct-memory fast paths. + fn compressible_plaintext() -> Vec { + b"rustfs-body-cache-e2e-regression-".repeat(20_000) + } + + /// backlog#1108: a raw data-movement read (decommission/rebalance copy) must + /// yield the STORED (compressed) representation, never the cached plaintext. + /// Serving the cache here writes decompressed bytes into the destination + /// pool under the compressed object's metadata — silent corruption. + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn raw_data_movement_read_serves_stored_bytes_not_cached_plaintext() { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let bucket = "e2e-body-cache-raw-move"; + let object = "compressed.bin"; + let plaintext = compressible_plaintext(); + let compressed = compress(&plaintext).await; + assert_ne!(compressed, plaintext, "fixture must be genuinely compressed"); + put_compressed_object(&set_disks, bucket, object, &plaintext, &compressed).await; + + // Prime the cache with the DECOMPRESSED plaintext. + let _guard = HookGuard::install(bucket, object, Bytes::from(plaintext.clone())); + + // decommission_object_migration_read_opts sets both raw_data_movement_read + // and data_movement; ReadPlan keys the stored-bytes path on the former. + let opts = ObjectOptions { + no_lock: true, + raw_data_movement_read: true, + ..Default::default() + }; + let (body, published_size) = read_to_vec(&set_disks, bucket, object, &opts).await; + + assert_eq!(body, compressed, "raw data-movement read must return the stored compressed bytes"); + assert_ne!(body, plaintext, "raw data-movement read must NOT return the cached plaintext"); + assert_eq!(published_size, compressed.len() as i64, "raw read must publish the stored size"); + } + + /// backlog#1109: on a cache hit for a compressed object the served + /// `object_info.size` must be the DECOMPRESSED length (what + /// ReadTransform::Compressed publishes), and the streamed body length must + /// match it. UploadPartCopy reads the copy length straight off this field. + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn compressed_cache_hit_publishes_decompressed_size() { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let bucket = "e2e-body-cache-compressed"; + let object = "compressed.bin"; + let plaintext = compressible_plaintext(); + let compressed = compress(&plaintext).await; + assert_ne!(compressed, plaintext, "fixture must be genuinely compressed"); + put_compressed_object(&set_disks, bucket, object, &plaintext, &compressed).await; + + let normal_opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + // Control: with no hook the real decode path must decompress to the + // plaintext and publish the decompressed length. This also proves the + // fixture round-trips (stored compressed bytes -> plaintext). + clear_get_object_body_cache_hook(); + let (control_body, control_size) = read_to_vec(&set_disks, bucket, object, &normal_opts).await; + assert_eq!(control_body, plaintext, "real read must decompress to the plaintext"); + assert_eq!(control_size, plaintext.len() as i64, "real read must publish the decompressed size"); + + // Hit: prime with the plaintext and read normally. The probe must serve + // it AND republish the decompressed length, not the stored size. + let _guard = HookGuard::install(bucket, object, Bytes::from(plaintext.clone())); + let (hit_body, hit_size) = read_to_vec(&set_disks, bucket, object, &normal_opts).await; + assert_eq!( + hit_size, + plaintext.len() as i64, + "cache hit must publish the decompressed size (backlog#1109)" + ); + assert_ne!(hit_size, compressed.len() as i64, "cache hit must not publish the stored compressed size"); + assert_eq!(hit_body.len() as i64, hit_size, "streamed length must equal the published size"); + assert_eq!(hit_body, plaintext, "cache hit must stream the plaintext"); + } + + /// backlog#1146: a restore read forces ReadPlan down the Plain branch, so a + /// compressed object yields its STORED bytes under the compressed size. The + /// cache (holding plaintext) must not be served. + #[tokio::test] + #[serial_test::serial(body_cache_hook)] + async fn restore_read_serves_stored_bytes_not_cached_plaintext() { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let bucket = "e2e-body-cache-restore"; + let object = "compressed.bin"; + let plaintext = compressible_plaintext(); + let compressed = compress(&plaintext).await; + assert_ne!(compressed, plaintext, "fixture must be genuinely compressed"); + put_compressed_object(&set_disks, bucket, object, &plaintext, &compressed).await; + + let _guard = HookGuard::install(bucket, object, Bytes::from(plaintext.clone())); + + let mut opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + opts.transition.restore_request.days = Some(1); + let (body, published_size) = read_to_vec(&set_disks, bucket, object, &opts).await; + + assert_eq!(body, compressed, "restore read must return the stored compressed bytes"); + assert_ne!(body, plaintext, "restore read must NOT return the cached plaintext"); + assert_eq!( + published_size, + compressed.len() as i64, + "restore read must publish the stored compressed size" + ); + } +} diff --git a/crates/obs/src/telemetry/dial9/disabled.rs b/crates/obs/src/telemetry/dial9/disabled.rs index de2b299e2..613028225 100644 --- a/crates/obs/src/telemetry/dial9/disabled.rs +++ b/crates/obs/src/telemetry/dial9/disabled.rs @@ -38,6 +38,13 @@ impl Dial9SessionGuard { } } +/// Mirrors the `Drop` the enabled guard uses to flush and seal the trace +/// segment, so a caller can drop the guard before `process::exit` under either +/// feature without `clippy::drop_non_drop` firing on this variant. +impl Drop for Dial9SessionGuard { + fn drop(&mut self) {} +} + /// Always fails: telemetry support is not compiled in. /// /// # Errors