Compare commits

..

2 Commits

Author SHA1 Message Date
overtrue bae602a5ef test(table-catalog): fold the admin object backend into the shared fake
Second half of PR2: the admin handler tests' TestTableCatalogObjectBackend (struct, inherent impl, trait impl, lock alias — 201 lines) is deleted and its instrumentation folded into the shared TestCatalogObjectBackend, which the admin tests now take through a type alias so no call site is renamed.

Two behavioral differences between the two fakes were found by the test suites rather than assumed away, and both are preserved:

- Lock observability: the admin fake implemented only acquire_write_lock, so the trait's default acquire_read_lock -> acquire_write_lock delegation made read acquisitions visible in lock_attempts. The shared fake implements both independently, so five fence/lock tests timed out until the read path also records attempts.

- Etag generation: the admin fake used content-addressed sha256 etags (its tests observe an etag and expect rewriting identical bytes to reproduce it) while the store fake uses an incrementing counter. Instead of silently picking one, the union carries a content_addressed_etags flag; the 80 admin construction sites go through TestCatalogObjectBackend::content_addressed() and the store tests keep counter semantics.

The six one-shot path-keyed injection knobs (fail/corrupt put, missing/fail read, put barrier) run before the store fake's attempt-indexed injection maps, matching each fake's original ordering.

Verification: cargo test -p rustfs --lib table_catalog 481 passed; --lib admin::handlers::table_catalog 167 passed; clippy --lib --tests -D warnings clean; make pre-commit green.

Ref rustfs/backlog#1837 (PR2, part 2).
2026-08-13 23:40:41 +08:00
overtrue d9071b1b83 test(table-catalog): move the store-side stateful object backend into test_support
First half of the issue's PR2: the store tests' TestCatalogObjectBackend cluster (state/record/locks/pause types, the seed/fail/pause instrumented inherent impl, the TableCatalogObjectBackend trait impl, and the BlockingObjectPublication/UnserializedTestPublication commit-publication fakes — 544 lines) moves verbatim from table_catalog/tests.rs into test_support.rs, with pub(crate) visibility on the items and fields the tests reach directly. Pure move, no behavior change; the admin handler tests' TestTableCatalogObjectBackend union (its put barrier / fail-path / lock-attempt instrumentation folding into this fake) is the second half.

Verification: cargo test -p rustfs --lib table_catalog 481 passed; clippy --lib --tests -D warnings clean; make pre-commit green.

Ref rustfs/backlog#1837 (PR2, part 1).
2026-08-13 20:17:28 +08:00
22 changed files with 1316 additions and 3712 deletions
+40 -194
View File
@@ -71,16 +71,10 @@ impl EncodedBlock {
const MODERN_MAX_TOTAL_SHARDS: usize = <reed_solomon_erasure::galois_8::Field as reed_solomon_erasure::Field>::ORDER;
const MODERN_REED_SOLOMON_CACHE_MAX_ENTRIES: usize = 64;
const LEGACY_REED_SOLOMON_CACHE_MAX_ENTRIES: usize = 16;
// Vec growth may retain twice the requested logical length. Keeping the logical
// workspace at half the budget bounds each cached workspace's shard allocation to 1 MiB.
const LEGACY_REED_SOLOMON_CACHE_MAX_LOGICAL_SHARD_BYTES_PER_WORKSPACE: usize = 512 * 1024;
type ModernReedSolomonCache = RwLock<HashMap<(usize, usize), Arc<ReedSolomon>>>;
type LegacyReedSolomonCache = RwLock<HashMap<(usize, usize), Arc<LegacyReedSolomonEncoder>>>;
static MODERN_REED_SOLOMON_CACHE: OnceLock<ModernReedSolomonCache> = OnceLock::new();
static LEGACY_REED_SOLOMON_CACHE: OnceLock<LegacyReedSolomonCache> = OnceLock::new();
/// Errors returned when constructing an [`Erasure`] codec.
#[derive(Debug, thiserror::Error)]
@@ -147,61 +141,43 @@ pub fn calc_shard_size_legacy(block_size: usize, data_shards: usize) -> usize {
struct LegacyReedSolomonEncoder {
data_shards: usize,
parity_shards: usize,
cache_workspaces: bool,
encoder_cache: RwLock<Option<reed_solomon_simd::ReedSolomonEncoder>>,
decoder_cache: RwLock<Option<reed_solomon_simd::ReedSolomonDecoder>>,
encoder_cache: std::sync::RwLock<Option<reed_solomon_simd::ReedSolomonEncoder>>,
decoder_cache: std::sync::RwLock<Option<reed_solomon_simd::ReedSolomonDecoder>>,
}
impl Clone for LegacyReedSolomonEncoder {
fn clone(&self) -> Self {
Self {
data_shards: self.data_shards,
parity_shards: self.parity_shards,
encoder_cache: std::sync::RwLock::new(None),
decoder_cache: std::sync::RwLock::new(None),
}
}
}
impl LegacyReedSolomonEncoder {
fn new(data_shards: usize, parity_shards: usize) -> io::Result<Self> {
Self::with_workspace_cache(data_shards, parity_shards, false)
}
fn with_workspace_cache(data_shards: usize, parity_shards: usize, cache_workspaces: bool) -> io::Result<Self> {
fn new(_data_shards: usize, _parity_shards: usize) -> io::Result<Self> {
Ok(Self {
data_shards,
parity_shards,
cache_workspaces,
encoder_cache: RwLock::new(None),
decoder_cache: RwLock::new(None),
data_shards: _data_shards,
parity_shards: _parity_shards,
encoder_cache: std::sync::RwLock::new(None),
decoder_cache: std::sync::RwLock::new(None),
})
}
fn logical_shard_bytes_upper_bound(&self, shard_len: usize) -> Option<usize> {
let aligned_shard_len = shard_len.checked_add(63)?.checked_div(64)?.checked_mul(64)?;
let high_rate_decoder_work_count = self
.parity_shards
.checked_next_power_of_two()?
.checked_add(self.data_shards)?
.checked_next_power_of_two()?;
let low_rate_decoder_work_count = self
.data_shards
.checked_next_power_of_two()?
.checked_add(self.parity_shards)?
.checked_next_power_of_two()?;
aligned_shard_len.checked_mul(high_rate_decoder_work_count.max(low_rate_decoder_work_count))
}
fn should_cache_workspace(&self, shard_len: usize) -> bool {
self.cache_workspaces
&& self
.logical_shard_bytes_upper_bound(shard_len)
.is_some_and(|bytes| bytes <= LEGACY_REED_SOLOMON_CACHE_MAX_LOGICAL_SHARD_BYTES_PER_WORKSPACE)
}
fn encode(&self, shards: SmallVec<[&mut [u8]; 16]>) -> io::Result<()> {
let mut shards_vec: Vec<&mut [u8]> = shards.into_vec();
if shards_vec.is_empty() {
return Ok(());
}
let shard_len = shards_vec[0].len();
let cached_encoder = self
.encoder_cache
.write()
.map_err(|_| io::Error::other("Failed to acquire encoder cache lock"))?
.take();
let mut encoder = {
match cached_encoder {
let mut cache_guard = self
.encoder_cache
.write()
.map_err(|_| io::Error::other("Failed to acquire encoder cache lock"))?;
match cache_guard.take() {
Some(mut cached) => {
if cached.reset(self.data_shards, self.parity_shards, shard_len).is_err() {
reed_solomon_simd::ReedSolomonEncoder::new(self.data_shards, self.parity_shards, shard_len)
@@ -228,15 +204,10 @@ impl LegacyReedSolomonEncoder {
}
}
drop(result);
if self.should_cache_workspace(shard_len) {
let mut cache = self
.encoder_cache
.write()
.map_err(|_| io::Error::other("Failed to return encoder to cache"))?;
if cache.is_none() {
*cache = Some(encoder);
}
}
*self
.encoder_cache
.write()
.map_err(|_| io::Error::other("Failed to return encoder to cache"))? = Some(encoder);
Ok(())
}
@@ -250,13 +221,13 @@ impl LegacyReedSolomonEncoder {
.find_map(|s| s.as_ref().map(|v| v.len()))
.ok_or_else(|| io::Error::other("No valid shards found for reconstruction"))?;
let cached_decoder = self
.decoder_cache
.write()
.map_err(|_| io::Error::other("Failed to acquire decoder cache lock"))?
.take();
let mut decoder = {
match cached_decoder {
let mut cache_guard = self
.decoder_cache
.write()
.map_err(|_| io::Error::other("Failed to acquire decoder cache lock"))?;
match cache_guard.take() {
Some(mut cached_decoder) => {
if let Err(e) = cached_decoder.reset(self.data_shards, self.parity_shards, shard_len) {
warn!("Failed to reset SIMD decoder: {:?}, creating new one", e);
@@ -303,15 +274,10 @@ impl LegacyReedSolomonEncoder {
drop(result);
if self.should_cache_workspace(shard_len) {
let mut cache = self
.decoder_cache
.write()
.map_err(|_| io::Error::other("Failed to return decoder to cache"))?;
if cache.is_none() {
*cache = Some(decoder);
}
}
*self
.decoder_cache
.write()
.map_err(|_| io::Error::other("Failed to return decoder to cache"))? = Some(decoder);
Ok(())
}
@@ -469,39 +435,6 @@ fn cached_modern_reed_solomon(data_shards: usize, parity_shards: usize) -> Resul
Ok(encoder)
}
fn cached_legacy_reed_solomon(data_shards: usize, parity_shards: usize) -> io::Result<Arc<LegacyReedSolomonEncoder>> {
let cache = LEGACY_REED_SOLOMON_CACHE.get_or_init(|| RwLock::new(HashMap::new()));
cached_legacy_reed_solomon_in(cache, data_shards, parity_shards)
}
fn cached_legacy_reed_solomon_in(
cache: &LegacyReedSolomonCache,
data_shards: usize,
parity_shards: usize,
) -> io::Result<Arc<LegacyReedSolomonEncoder>> {
let key = (data_shards, parity_shards);
if let Some(encoder) = cache
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(&key)
.cloned()
{
return Ok(encoder);
}
let mut cache = cache.write().unwrap_or_else(|poisoned| poisoned.into_inner());
if let Some(existing) = cache.get(&key) {
return Ok(Arc::clone(existing));
}
if cache.len() < LEGACY_REED_SOLOMON_CACHE_MAX_ENTRIES {
let encoder = Arc::new(LegacyReedSolomonEncoder::with_workspace_cache(data_shards, parity_shards, true)?);
cache.insert(key, Arc::clone(&encoder));
return Ok(encoder);
}
drop(cache);
Ok(Arc::new(LegacyReedSolomonEncoder::new(data_shards, parity_shards)?))
}
fn encode_parity_shards<F>(shards: &mut [Option<Vec<u8>>], data_shards: usize, parity_shards: usize, encode: F) -> io::Result<()>
where
F: FnOnce(SmallVec<[&mut [u8]; 16]>) -> io::Result<()>,
@@ -618,7 +551,7 @@ pub struct Erasure {
pub data_shards: usize,
pub parity_shards: usize,
encoder: Option<ReedSolomonEncoder>,
legacy_encoder: Option<Arc<LegacyReedSolomonEncoder>>,
legacy_encoder: Option<LegacyReedSolomonEncoder>,
pub block_size: usize,
uses_legacy: bool,
_id: Uuid,
@@ -754,7 +687,7 @@ impl Erasure {
let legacy_encoder = if uses_legacy && parity_shards > 0 {
Some(
cached_legacy_reed_solomon(data_shards, parity_shards)
LegacyReedSolomonEncoder::new(data_shards, parity_shards)
.map_err(|source| ErasureConstructionError::LegacyEncoder { source })?,
)
} else {
@@ -1472,7 +1405,7 @@ mod tests {
assert_eq!(cloned.block_size, legacy.block_size);
assert!(cloned.uses_legacy);
let data = b"legacy clone should preserve SIMD codec behavior";
let data = b"legacy clone should keep independent SIMD caches";
let encoded = cloned.encode_data(data).expect("legacy clone should encode");
let mut shards = optional_shards(&encoded);
shards[0] = None;
@@ -1480,93 +1413,6 @@ mod tests {
assert_eq!(recover_data(&shards, cloned.data_shards, data.len()), data);
}
#[test]
fn legacy_codecs_share_process_cache_across_erasure_instances() {
let first = Erasure::new_with_options(6, 3, 64, true)
.legacy_encoder
.expect("legacy codec should be initialized");
let second = Erasure::new_with_options(6, 3, 128, true)
.legacy_encoder
.expect("same legacy shard layout should be initialized");
assert!(Arc::ptr_eq(&first, &second));
}
#[test]
fn legacy_workspace_cache_rejects_oversize_buffers_and_isolates_layouts() {
let four_plus_two = Erasure::new_with_options(4, 2, 64, true)
.legacy_encoder
.expect("legacy codec should be initialized");
let four_plus_one = Erasure::new_with_options(4, 1, 64, true)
.legacy_encoder
.expect("distinct parity layout should be initialized");
let three_plus_two = Erasure::new_with_options(3, 2, 64, true)
.legacy_encoder
.expect("distinct data layout should be initialized");
assert!(!Arc::ptr_eq(&four_plus_two, &four_plus_one));
assert!(!Arc::ptr_eq(&four_plus_two, &three_plus_two));
assert_eq!(four_plus_two.logical_shard_bytes_upper_bound(64 * 1024), Some(512 * 1024));
assert!(four_plus_two.should_cache_workspace(64 * 1024));
assert!(!four_plus_two.should_cache_workspace(64 * 1024 + 1));
let nine_plus_seven =
LegacyReedSolomonEncoder::with_workspace_cache(9, 7, true).expect("9+7 legacy codec should construct");
assert_eq!(nine_plus_seven.logical_shard_bytes_upper_bound(16 * 1024), Some(512 * 1024));
assert!(nine_plus_seven.should_cache_workspace(16 * 1024));
assert!(!nine_plus_seven.should_cache_workspace(16 * 1024 + 1));
let uncached = LegacyReedSolomonEncoder::new(4, 2).expect("uncached legacy codec should construct");
assert!(!uncached.should_cache_workspace(64));
}
#[test]
fn saturated_legacy_codec_cache_does_not_retain_more_workspaces() {
let cache = RwLock::new(HashMap::new());
for parity_shards in 1..=LEGACY_REED_SOLOMON_CACHE_MAX_ENTRIES {
let cached =
cached_legacy_reed_solomon_in(&cache, 32, parity_shards).expect("cacheable legacy codec should construct");
assert!(cached.cache_workspaces);
}
let uncached =
cached_legacy_reed_solomon_in(&cache, 31, 1).expect("uncached legacy codec should construct after saturation");
assert!(!uncached.cache_workspaces);
assert_eq!(
cache.read().expect("cache lock should remain healthy").len(),
LEGACY_REED_SOLOMON_CACHE_MAX_ENTRIES
);
}
#[test]
fn concurrent_legacy_codecs_preserve_byte_exact_results() {
let barrier = Arc::new(std::sync::Barrier::new(2));
let payloads = [vec![0x35; 257], vec![0xca; 1025]];
std::thread::scope(|scope| {
let handles = payloads.each_ref().map(|payload| {
let barrier = Arc::clone(&barrier);
scope.spawn(move || {
let erasure = Erasure::new_with_options(6, 3, 2048, true);
barrier.wait();
let encoded = erasure.encode_data(payload).expect("concurrent legacy encode should succeed");
barrier.wait();
let mut shards = optional_shards(&encoded);
shards[0] = None;
erasure
.decode_data(&mut shards)
.expect("concurrent legacy decode should reconstruct the missing shard");
recover_data(&shards, erasure.data_shards, payload.len())
})
});
for (handle, payload) in handles.into_iter().zip(payloads.iter()) {
assert_eq!(handle.join().expect("concurrent legacy codec worker should not panic"), *payload);
}
});
}
#[test]
fn legacy_verify_reports_invalid_empty_valid_and_corrupt_parity_sets() {
let legacy = LegacyReedSolomonEncoder::new(2, 2).expect("legacy encoder should construct");
+8 -77
View File
@@ -15,7 +15,6 @@
use super::*;
use crate::io_support::rio::Index;
use std::mem::MaybeUninit;
#[cfg(feature = "rio-v2")]
const DARE_PAYLOAD_SIZE: i64 = 64 * 1024;
@@ -923,7 +922,7 @@ struct SkipReader<R> {
inner: R,
bytes_to_skip: usize,
bytes_skipped: usize,
scratch: Box<[MaybeUninit<u8>]>,
scratch: Vec<u8>,
}
impl<R: AsyncRead + Unpin + Send + Sync> SkipReader<R> {
@@ -932,7 +931,7 @@ impl<R: AsyncRead + Unpin + Send + Sync> SkipReader<R> {
inner,
bytes_to_skip,
bytes_skipped: 0,
scratch: Box::<[u8]>::new_uninit_slice(8192),
scratch: vec![0u8; 8192],
}
}
}
@@ -944,7 +943,7 @@ impl<R: AsyncRead + Unpin + Send + Sync> AsyncRead for SkipReader<R> {
while this.bytes_skipped < this.bytes_to_skip {
let remaining = this.bytes_to_skip - this.bytes_skipped;
let scratch_len = remaining.min(this.scratch.len());
let mut scratch_buf = ReadBuf::uninit(&mut this.scratch[..scratch_len]);
let mut scratch_buf = ReadBuf::new(&mut this.scratch[..scratch_len]);
match Pin::new(&mut this.inner).poll_read(cx, &mut scratch_buf) {
Poll::Pending => return Poll::Pending,
Poll::Ready(Err(err)) => return Poll::Ready(Err(err)),
@@ -975,7 +974,7 @@ pub struct RangedDecompressReader<R: AsyncRead + Unpin + Send + Sync + 'static>
target_length: usize,
current_offset: usize,
bytes_returned: usize,
scratch: Box<[MaybeUninit<u8>]>,
scratch: Vec<u8>,
drain_on_done: bool,
drain_task: Option<tokio::task::JoinHandle<()>>,
}
@@ -1013,7 +1012,7 @@ impl<R: AsyncRead + Unpin + Send + Sync + 'static> RangedDecompressReader<R> {
target_length: actual_length,
current_offset: 0,
bytes_returned: 0,
scratch: Box::<[u8]>::new_uninit_slice(8192),
scratch: vec![0u8; 8192],
drain_on_done,
drain_task: None,
})
@@ -1063,7 +1062,7 @@ impl<R: AsyncRead + Unpin + Send + Sync + 'static> AsyncRead for RangedDecompres
}
let scratch_len = std::cmp::min(this.scratch.len(), std::cmp::max(buf_capacity, 1));
let mut temp_read_buf = ReadBuf::uninit(&mut this.scratch[..scratch_len]);
let mut temp_read_buf = ReadBuf::new(&mut this.scratch[..scratch_len]);
let Some(inner) = this.inner.as_mut() else {
return Poll::Ready(Ok(()));
@@ -1115,8 +1114,7 @@ impl<R: AsyncRead + Unpin + Send + Sync + 'static> AsyncRead for RangedDecompres
);
if bytes_to_return > 0 {
let data_slice =
&temp_read_buf.filled()[data_start_in_buffer..data_start_in_buffer + bytes_to_return];
let data_slice = &this.scratch[data_start_in_buffer..data_start_in_buffer + bytes_to_return];
buf.put_slice(data_slice);
this.bytes_returned += bytes_to_return;
@@ -1135,7 +1133,7 @@ impl<R: AsyncRead + Unpin + Send + Sync + 'static> AsyncRead for RangedDecompres
std::cmp::min(n, std::cmp::min(buf.remaining(), this.target_length - this.bytes_returned));
if bytes_to_return > 0 {
buf.put_slice(&temp_read_buf.filled()[..bytes_to_return]);
buf.put_slice(&this.scratch[..bytes_to_return]);
this.bytes_returned += bytes_to_return;
tracing::trace!("Returned {} bytes at offset {}", bytes_to_return, old_offset);
@@ -1265,43 +1263,6 @@ mod tests {
use temp_env::async_with_vars;
use tokio::io::AsyncReadExt;
#[derive(Debug)]
struct PendingPartialReader {
data: &'static [u8],
position: usize,
pending: bool,
}
impl PendingPartialReader {
fn new(data: &'static [u8]) -> Self {
Self {
data,
position: 0,
pending: true,
}
}
}
impl AsyncRead for PendingPartialReader {
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
if self.pending {
self.pending = false;
cx.waker().wake_by_ref();
return Poll::Pending;
}
if self.position == self.data.len() {
return Poll::Ready(Ok(()));
}
let length = buf.remaining().min(3).min(self.data.len() - self.position);
let end = self.position + length;
buf.put_slice(&self.data[self.position..end]);
self.position = end;
self.pending = true;
Poll::Ready(Ok(()))
}
}
const TEST_DIRECT_KEY_HEADER: &str = "x-rustfs-test-direct-key";
const TEST_OBJECT_KEY_HEADER: &str = "x-rustfs-test-object-key";
const TEST_NONCE_HEADER: &str = "x-rustfs-test-nonce";
@@ -1439,36 +1400,6 @@ mod tests {
assert_eq!(result, b"World");
}
#[tokio::test]
async fn uninitialized_scratch_preserves_partial_pending_and_eof_reads() {
let mut skipped = SkipReader::new(PendingPartialReader::new(b"0123456789abcdef"), 5);
let mut skipped_output = Vec::new();
skipped
.read_to_end(&mut skipped_output)
.await
.expect("skip reader should survive partial pending reads through EOF");
assert_eq!(skipped_output, b"56789abcdef");
let mut ranged = RangedDecompressReader::new(PendingPartialReader::new(b"0123456789abcdef"), 5, 7, 16)
.expect("valid range should construct");
let mut ranged_output = Vec::new();
ranged
.read_to_end(&mut ranged_output)
.await
.expect("range reader should survive partial pending reads through EOF");
assert_eq!(ranged_output, b"56789ab");
}
#[tokio::test]
async fn uninitialized_skip_scratch_reports_early_eof() {
let mut reader = SkipReader::new(PendingPartialReader::new(b"short"), 6);
let error = reader
.read_to_end(&mut Vec::new())
.await
.expect_err("EOF before the skip boundary must remain visible");
assert_eq!(error.kind(), std::io::ErrorKind::UnexpectedEof);
}
#[tokio::test]
async fn test_ranged_decompress_reader_from_start() {
let original_data = b"Hello, World! This is a test.";
File diff suppressed because it is too large Load Diff
+29 -546
View File
@@ -880,44 +880,12 @@ mod prepared_get_object_metadata_tests {
use super::*;
use crate::ecstore_validation_blackbox::make_local_set_disks;
use crate::object_api::{BLOCK_SIZE_V2, PutObjReader};
use crate::set_disk::core::io_primitives::{bounded_metadata_fanout_order, disk_call_counters, rename_fanout_barrier};
use crate::set_disk::core::io_primitives::disk_call_counters;
use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions};
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
use crate::test_metrics::CapturingRecorder;
use http::HeaderMap;
use tokio::io::AsyncReadExt;
const READ_VERSION_BARRIER_GUARD: std::time::Duration = std::time::Duration::from_secs(10);
fn object_with_initial_data_shards(bucket: &str, prefix: &str) -> String {
(0..1000)
.map(|index| format!("{prefix}-{index}.bin"))
.find(|name| {
let order = bounded_metadata_fanout_order(bucket, name, 4, 2);
let distribution = FileInfo::new(&[bucket, name].join("/"), 2, 2).erasure.distribution;
let mut seen = [false; 2];
for disk_index in order.into_iter().take(3) {
if let Some(block_index @ 1..=2) = distribution.get(disk_index).copied() {
seen[block_index - 1] = true;
}
}
seen.into_iter().all(|seen| seen)
})
.expect("test should find an object whose initial fanout covers both data shards")
}
fn bounded_spare_disk_index(bucket: &str, object: &str) -> usize {
*bounded_metadata_fanout_order(bucket, object, 4, 2)
.get(3)
.expect("4-disk test geometry should leave one bounded spare disk")
}
fn bounded_slow_initial_disk_index(bucket: &str, object: &str) -> usize {
*bounded_metadata_fanout_order(bucket, object, 4, 2)
.get(2)
.expect("4-disk test geometry should include a third initial metadata disk")
}
#[tokio::test]
async fn prepared_metadata_is_consumed_exactly_once() {
let snapshot = GetObjectFileInfo::owned(FileInfo::default(), Vec::new(), Vec::new());
@@ -1034,307 +1002,6 @@ mod prepared_get_object_metadata_tests {
);
}
#[test]
#[serial_test::serial(body_cache_hook)]
fn inline_data_read_early_stop_reader_returns_exact_body() {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime should build");
let bucket = "inline-data-read-early-stop-reader";
let object = object_with_initial_data_shards(bucket, "inline-data-read-early-stop-reader-object");
let payload = b"inline early-stop reader payload".repeat(256);
let recorder = CapturingRecorder::default();
let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled();
rustfs_io_metrics::set_get_stage_metrics_enabled(true);
let (restored, object_size, calls_total) = metrics::with_local_recorder(&recorder, || {
runtime.block_on(async {
let (_dirs, set_disks) = make_local_set_disks(4, 2).await;
let opts = ObjectOptions {
no_lock: true,
..Default::default()
};
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut put_reader = PutObjReader::from_vec(payload.clone());
set_disks
.put_object(bucket, &object, &mut put_reader, &opts)
.await
.expect("inline object should be written");
temp_env::async_with_vars(
[
("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")),
],
async {
let slow_initial_disk = bounded_slow_initial_disk_index(bucket, &object);
let barrier =
rename_fanout_barrier::arm(&object, slow_initial_disk, rename_fanout_barrier::PHASE_READ_VERSION);
let calls = disk_call_counters::observe(&object);
let set_disks_for_read = Arc::clone(&set_disks);
let opts_for_read = opts.clone();
let object_for_read = object.clone();
let mut open_reader = tokio::spawn(async move {
set_disks_for_read
.get_object_reader(bucket, &object_for_read, None, HeaderMap::new(), &opts_for_read)
.await
});
tokio::time::timeout(READ_VERSION_BARRIER_GUARD, barrier.wait_until_paused())
.await
.expect("bounded inline GET should pause a slow initial metadata read");
let mut reader = tokio::time::timeout(READ_VERSION_BARRIER_GUARD, &mut open_reader)
.await
.expect("production inline GET should return before the paused metadata response")
.expect("inline GET reader task should not panic")
.expect("inline GET reader should open");
let object_size = reader.object_info.size;
let mut restored = Vec::new();
reader
.stream
.read_to_end(&mut restored)
.await
.expect("inline GET body should stream");
(restored, object_size, calls.total(disk_call_counters::KIND_READ_VERSION))
},
)
.await
})
});
rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate);
assert_eq!(object_size, payload.len() as i64);
assert_eq!(restored, payload);
assert_eq!(calls_total, 4, "bounded production GET should schedule the initial quorum plus one spare");
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_scheduled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![4.0],
"bounded production GET should record all scheduled metadata tasks"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_completed",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![3.0],
"bounded production GET should record only observed metadata responses as completed"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_cancelled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![1.0],
"bounded production GET should record the aborted slow metadata task"
);
}
#[test]
#[serial_test::serial(body_cache_hook)]
fn prepared_metadata_uses_full_fanout_even_when_data_read_early_stop_is_enabled() {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime should build");
let bucket = "prepared-metadata-early-stop-enabled";
let object = object_with_initial_data_shards(bucket, "prepared-metadata-early-stop-enabled-object");
let payload = b"prepared metadata early-stop enabled payload".repeat(16);
let recorder = CapturingRecorder::default();
let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled();
rustfs_io_metrics::set_get_stage_metrics_enabled(true);
let (restored, calls_total) = metrics::with_local_recorder(&recorder, || {
runtime.block_on(async {
let (_dirs, set_disks) = make_local_set_disks(4, 2).await;
let opts = ObjectOptions {
no_lock: true,
..Default::default()
};
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut put_reader = PutObjReader::from_vec(payload.clone());
set_disks
.put_object(bucket, &object, &mut put_reader, &opts)
.await
.expect("object should be written");
temp_env::async_with_vars(
[
("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")),
],
async {
let calls = disk_call_counters::observe(&object);
let metadata = set_disks
.prepare_get_object_metadata(bucket, &object, &opts)
.await
.expect("prepared metadata should resolve");
let calls_total = calls.total(disk_call_counters::KIND_READ_VERSION);
let mut reader = set_disks
.get_object_reader_with_prepared_metadata(bucket, &object, None, HeaderMap::new(), &opts, metadata)
.await
.expect("prepared body reader should open");
let mut restored = Vec::new();
reader
.stream
.read_to_end(&mut restored)
.await
.expect("prepared body should stream");
(restored, calls_total)
},
)
.await
})
});
rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate);
assert_eq!(restored, payload);
assert_eq!(
calls_total, 4,
"prepared metadata must opt out of data-read early-stop until the read shape is known"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_scheduled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![4.0],
"prepared metadata should schedule the full metadata fanout"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_completed",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![4.0],
"prepared metadata must wait for every scheduled metadata response"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_cancelled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![0.0],
"prepared metadata must not cancel metadata responses"
);
}
#[test]
#[serial_test::serial(body_cache_hook)]
fn data_read_early_stop_request_shapes_full_wait_in_production_reader() {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime should build");
let bucket = "data-read-early-stop-shape-reader";
let payload = b"shape-gated inline reader payload".repeat(256);
for (object_prefix, range, configure_opts, expected_body) in [
(
"data-read-early-stop-range-reader-object",
Some(HTTPRangeSpec {
start: 0,
end: 3,
is_suffix_length: false,
}),
None,
payload[..4].to_vec(),
),
("data-read-early-stop-part-reader-object", None, Some(1), payload.clone()),
] {
let recorder = CapturingRecorder::default();
let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled();
rustfs_io_metrics::set_get_stage_metrics_enabled(true);
let (restored, calls_total) = metrics::with_local_recorder(&recorder, || {
runtime.block_on(async {
let (_dirs, set_disks) = make_local_set_disks(4, 2).await;
let object = object_with_initial_data_shards(bucket, object_prefix);
let mut opts = ObjectOptions {
no_lock: true,
..Default::default()
};
opts.part_number = configure_opts;
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut put_reader = PutObjReader::from_vec(payload.clone());
set_disks
.put_object(bucket, &object, &mut put_reader, &opts)
.await
.expect("inline object should be written");
temp_env::async_with_vars(
[
("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")),
],
async {
let calls = disk_call_counters::observe(&object);
let mut reader = set_disks
.get_object_reader(bucket, &object, range, HeaderMap::new(), &opts)
.await
.expect("shape-gated GET reader should open");
let mut restored = Vec::new();
reader
.stream
.read_to_end(&mut restored)
.await
.expect("shape-gated GET body should stream");
(restored, calls.total(disk_call_counters::KIND_READ_VERSION))
},
)
.await
})
});
rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate);
assert_eq!(restored, expected_body);
assert_eq!(calls_total, 4, "shape-gated production GET should keep full metadata fanout");
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_scheduled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![4.0],
"shape-gated production GET should schedule the full metadata fanout"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_completed",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![4.0],
"shape-gated production GET must wait for every scheduled metadata response"
);
assert_eq!(
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_cancelled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)]
),
vec![0.0],
"shape-gated production GET must not cancel metadata responses"
);
}
}
#[tokio::test]
#[serial_test::serial(body_cache_hook)]
async fn prepared_reader_rebuilds_object_info_when_precomputed_value_is_absent() {
@@ -1438,7 +1105,7 @@ impl SetDisks {
object: &str,
opts: &ObjectOptions,
) -> Result<PreparedGetObjectMetadata> {
let snapshot = self.get_object_fileinfo(bucket, object, opts, true, false).await?;
let snapshot = self.get_object_fileinfo(bucket, object, opts, true, true).await?;
let object_info = build_get_object_info(snapshot.fi(), bucket, object, opts.versioned || opts.version_suspended);
Ok(PreparedGetObjectMetadata {
snapshot,
@@ -2792,9 +2459,6 @@ pub struct SetDisks {
get_object_metadata_cache: moka::future::Cache<GetObjectMetadataCacheKey, Arc<GetObjectMetadataCacheEntry>>,
get_object_metadata_cache_hash_builder: std::collections::hash_map::RandomState,
get_object_metadata_cache_generations: Arc<[AtomicU64]>,
/// GET codecs keyed by every persisted layout dimension that affects
/// decoding. Clones of a set share the memoized shells.
erasure_cache: Arc<ErasureCache>,
pub lockers: Vec<Arc<dyn LockClient>>,
shared_lockers: Arc<[Arc<dyn LockClient>]>,
local_lock_manager: Arc<rustfs_lock::GlobalLockManager>,
@@ -2817,137 +2481,6 @@ pub struct SetDisks {
storage_class_config_override: Arc<std::sync::RwLock<Option<Arc<storageclass::Config>>>>,
}
const ERASURE_CACHE_MAX_ENTRIES: usize = 32;
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
struct ErasureCacheKey {
data_shards: usize,
parity_shards: usize,
block_size: usize,
uses_legacy: bool,
}
struct ErasureCache {
entries: parking_lot::RwLock<HashMap<ErasureCacheKey, Arc<coding::Erasure>>>,
}
impl std::fmt::Debug for ErasureCache {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("ErasureCache")
.field("entries", &self.entries.read().len())
.finish()
}
}
impl ErasureCache {
fn new() -> Self {
Self {
entries: parking_lot::RwLock::new(HashMap::new()),
}
}
fn get_or_try_insert(
&self,
key: ErasureCacheKey,
) -> std::result::Result<Arc<coding::Erasure>, coding::ErasureConstructionError> {
if let Some(erasure) = self.entries.read().get(&key) {
return Ok(Arc::clone(erasure));
}
// Serialize first construction for a key so concurrent cold GETs still
// create exactly one shell. Codec construction never awaits.
let mut entries = self.entries.write();
if let Some(erasure) = entries.get(&key) {
return Ok(Arc::clone(erasure));
}
let erasure = Arc::new(coding::Erasure::try_new_with_options(
key.data_shards,
key.parity_shards,
key.block_size,
key.uses_legacy,
)?);
if entries.len() < ERASURE_CACHE_MAX_ENTRIES {
entries.insert(key, Arc::clone(&erasure));
}
Ok(erasure)
}
fn get_for_file_info(&self, fi: &FileInfo) -> Result<Arc<coding::Erasure>> {
self.get_or_try_insert(ErasureCacheKey {
data_shards: fi.erasure.data_blocks,
parity_shards: fi.erasure.parity_blocks,
block_size: fi.erasure.block_size,
uses_legacy: fi.uses_legacy_checksum,
})
.map_err(Error::from)
}
}
#[cfg(test)]
mod erasure_cache_tests {
use super::*;
#[test]
fn reuses_shells_and_keeps_every_layout_dimension_in_the_key() {
let cache = ErasureCache::new();
let base = ErasureCacheKey {
data_shards: 4,
parity_shards: 2,
block_size: 1_048_576,
uses_legacy: false,
};
let first = cache.get_or_try_insert(base).expect("modern shell should construct");
let reused = cache.get_or_try_insert(base).expect("same modern shell should be cached");
assert!(Arc::ptr_eq(&first, &reused));
for distinct in [
ErasureCacheKey { data_shards: 3, ..base },
ErasureCacheKey {
parity_shards: 1,
..base
},
ErasureCacheKey {
block_size: 524_288,
..base
},
ErasureCacheKey {
uses_legacy: true,
..base
},
] {
let shell = cache.get_or_try_insert(distinct).expect("distinct shell should construct");
assert!(!Arc::ptr_eq(&first, &shell));
}
assert_eq!(cache.entries.read().len(), 5);
}
#[test]
fn does_not_cache_invalid_layouts_or_grow_past_the_bound() {
let cache = ErasureCache::new();
let invalid = ErasureCacheKey {
data_shards: 4,
parity_shards: 2,
block_size: 0,
uses_legacy: false,
};
assert!(cache.get_or_try_insert(invalid).is_err());
assert!(cache.entries.read().is_empty());
for block_size in 1..=(ERASURE_CACHE_MAX_ENTRIES + 1) {
cache
.get_or_try_insert(ErasureCacheKey {
data_shards: 4,
parity_shards: 2,
block_size,
uses_legacy: false,
})
.expect("bounded cache fixture should construct");
}
assert_eq!(cache.entries.read().len(), ERASURE_CACHE_MAX_ENTRIES);
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
struct GetObjectMetadataCacheKey {
bucket: Arc<str>,
@@ -3346,7 +2879,6 @@ impl SetDisks {
.map(|_| AtomicU64::new(0))
.collect::<Vec<_>>(),
),
erasure_cache: Arc::new(ErasureCache::new()),
lockers,
shared_lockers,
// Sourced from the instance context so each instance owns its lock
@@ -3879,15 +3411,9 @@ fn collect_inline_data_shard_fileinfos_by_index<'a>(
if block_index == 0 || block_index > data_shards {
continue;
}
if file_info.erasure.index != block_index {
continue;
}
if !file_info.has_valid_erasure_geometry() {
continue;
}
if !core::io_primitives::metadata_early_stop_candidate_matches(file_info, fi) {
continue;
}
if file_info.data.as_ref().is_none_or(|data| data.is_empty()) {
continue;
}
@@ -9695,9 +9221,6 @@ mod tests {
HashAlgorithm::HighwayHash256S
};
let shards = erasure.encode_data(payload).expect("payload should encode");
let version_id = Some(Uuid::new_v4());
let data_dir = Some(Uuid::new_v4());
let mod_time = Some(OffsetDateTime::now_utc());
let mut files = Vec::with_capacity(shards.len());
for shard in shards {
@@ -9710,16 +9233,6 @@ mod tests {
writer.shutdown().await.expect("inline writer should shutdown");
let data = writer.into_inline_data().expect("inline data should be retained");
let mut file = FileInfo::new("bucket/object", erasure.data_shards, erasure.parity_shards);
file.volume = "bucket".to_string();
file.name = "object".to_string();
file.size = i64::try_from(payload.len()).expect("test payload should fit i64");
file.is_latest = true;
file.version_id = version_id;
file.data_dir = data_dir;
file.mod_time = mod_time;
file.metadata.insert("etag".to_string(), "etag-inline".to_string());
file.add_object_part(1, "part-etag-inline".to_string(), payload.len(), file.mod_time, file.size, None, None);
file.set_inline_data();
file.erasure.index = files.len() + 1;
file.data = Some(Bytes::from(data));
files.push(file);
@@ -9732,40 +9245,16 @@ mod tests {
inline_bitrot_files_for_payload_with_mode(payload, false).await
}
fn disk_ordered_fileinfos(files: &[FileInfo]) -> Vec<FileInfo> {
let distribution = &files
.first()
.expect("inline data shard fixture should include metadata")
.erasure
.distribution;
distribution
.iter()
.map(|block_index| {
files
.get(block_index.checked_sub(1).expect("erasure block indexes are one-based"))
.expect("inline data shard fixture should include every distributed shard")
.clone()
})
.collect()
}
fn inline_data_shard_fileinfo(
name: &str,
data_blocks: usize,
parity_blocks: usize,
erasure_index: usize,
distribution: &[usize],
data: Option<&'static [u8]>,
) -> FileInfo {
let mut fi = FileInfo::new("object", data_blocks, parity_blocks);
fi.name = "object".to_string();
fi.volume = "bucket".to_string();
fi.size = 4;
fi.is_latest = true;
fi.data_dir = Some(Uuid::nil());
fi.mod_time = Some(OffsetDateTime::UNIX_EPOCH);
fi.metadata.insert("etag".to_string(), "etag-inline".to_string());
fi.add_object_part(1, "part-etag-inline".to_string(), 4, fi.mod_time, 4, None, None);
fi.set_inline_data();
let mut fi = FileInfo::new(name, data_blocks, parity_blocks);
fi.name = name.to_string();
fi.erasure.index = erasure_index;
fi.erasure.distribution = distribution.to_vec();
fi.data = data.map(Bytes::from_static);
@@ -9775,41 +9264,36 @@ mod tests {
#[test]
fn collect_inline_data_shards_by_index_uses_distribution_order() {
let distribution = vec![3, 1, 5, 2, 4, 6];
let mut fi = inline_data_shard_fileinfo(4, 2, 1, &distribution, Some(b"x"));
fi.erasure.index = 1;
let mut fi = FileInfo::new("object", 4, 2);
fi.erasure.distribution = distribution.clone();
let files = vec![
inline_data_shard_fileinfo(4, 2, 3, &distribution, Some(b"c")),
inline_data_shard_fileinfo(4, 2, 1, &distribution, Some(b"a")),
inline_data_shard_fileinfo(4, 2, 5, &distribution, Some(b"p")),
inline_data_shard_fileinfo(4, 2, 2, &distribution, Some(b"b")),
inline_data_shard_fileinfo(4, 2, 4, &distribution, Some(b"d")),
inline_data_shard_fileinfo(4, 2, 6, &distribution, Some(b"q")),
inline_data_shard_fileinfo("block-3", 4, 2, 3, &distribution, Some(b"c")),
inline_data_shard_fileinfo("block-1", 4, 2, 1, &distribution, Some(b"a")),
inline_data_shard_fileinfo("parity-5", 4, 2, 5, &distribution, Some(b"p")),
inline_data_shard_fileinfo("block-2", 4, 2, 2, &distribution, Some(b"b")),
inline_data_shard_fileinfo("block-4", 4, 2, 4, &distribution, Some(b"d")),
inline_data_shard_fileinfo("parity-6", 4, 2, 6, &distribution, Some(b"q")),
];
let data_files =
collect_inline_data_shard_fileinfos_by_index(&files, &fi, 4, |_| true).expect("all data shards should be collected");
assert_eq!(
data_files
.iter()
.map(|file| file.data.as_deref().expect("fixture carries inline bytes"))
.collect::<Vec<_>>(),
[b"a".as_slice(), b"b".as_slice(), b"c".as_slice(), b"d".as_slice()]
data_files.iter().map(|file| file.name.as_str()).collect::<Vec<_>>(),
["block-1", "block-2", "block-3", "block-4"]
);
}
#[test]
fn collect_inline_data_shards_by_index_rejects_missing_data_shard() {
let distribution = vec![1, 2, 3, 4];
let mut fi = inline_data_shard_fileinfo(2, 2, 1, &distribution, Some(b"x"));
fi.erasure.index = 1;
let mut fi = FileInfo::new("object", 2, 2);
fi.erasure.distribution = distribution.clone();
let files = vec![
inline_data_shard_fileinfo(2, 2, 1, &distribution, Some(b"a")),
inline_data_shard_fileinfo(2, 2, 2, &distribution, None),
inline_data_shard_fileinfo(2, 2, 3, &distribution, Some(b"p")),
inline_data_shard_fileinfo(2, 2, 4, &distribution, Some(b"q")),
inline_data_shard_fileinfo("block-1", 2, 2, 1, &distribution, Some(b"a")),
inline_data_shard_fileinfo("block-2", 2, 2, 2, &distribution, None),
inline_data_shard_fileinfo("parity-3", 2, 2, 3, &distribution, Some(b"p")),
inline_data_shard_fileinfo("parity-4", 2, 2, 4, &distribution, Some(b"q")),
];
assert!(collect_inline_data_shard_fileinfos_by_index(&files, &fi, 2, |_| true).is_none());
@@ -9942,8 +9426,10 @@ mod tests {
let payload = vec![b'i'; 192 * 1024];
let (erasure, files, _read_length, _checksum_algo) = inline_bitrot_files_for_payload(&payload).await;
let fi = files[0].clone();
let disk_files = disk_ordered_fileinfos(&files);
let mut fi = FileInfo::new("bucket/object", erasure.data_shards, erasure.parity_shards);
fi.size = payload.len() as i64;
fi.data = files[0].data.clone();
fi.add_object_part(1, String::new(), payload.len(), None, payload.len() as i64, None, None);
let disks = vec![Some(disk); erasure.total_shard_count()];
let metrics_size_bucket = rustfs_io_metrics::get_object_size_bucket(fi.size);
@@ -9951,9 +9437,8 @@ mod tests {
let body = SetDisks::try_get_object_direct_data_shards_with_fileinfo(
"bucket",
"object",
Arc::new(ErasureCache::new()),
&fi,
&disk_files,
&files,
&disks,
true,
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART,
@@ -9984,8 +9469,10 @@ mod tests {
let payload = vec![b'v'; 64 * 1024];
let payload_size = i64::try_from(payload.len()).expect("test payload size should fit i64");
let (erasure, files, _read_length, _checksum_algo) = inline_bitrot_files_for_payload(&payload).await;
let fi = files[0].clone();
let disk_files = disk_ordered_fileinfos(&files);
let mut fi = FileInfo::new("bucket/object", erasure.data_shards, erasure.parity_shards);
fi.size = payload_size;
fi.data = files[0].data.clone();
fi.add_object_part(1, String::new(), payload.len(), None, payload_size, None, None);
let mut object_info = ObjectInfo {
size: payload_size,
@@ -10015,9 +9502,8 @@ mod tests {
let body = SetDisks::try_get_object_direct_data_shards_with_fileinfo(
"bucket",
"object",
Arc::new(ErasureCache::new()),
&fi,
&disk_files,
&files,
&vec![Some(disk); erasure.total_shard_count()],
true,
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART,
@@ -10096,7 +9582,6 @@ mod tests {
let body = SetDisks::try_get_object_direct_data_shards_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
&fi,
&files,
&disks,
@@ -10182,7 +9667,6 @@ mod tests {
SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
range_offset,
range_length as i64,
&mut writer,
@@ -10294,7 +9778,6 @@ mod tests {
SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
0,
total_size as i64,
&mut writer,
+125 -210
View File
@@ -2296,157 +2296,140 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
}
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
let commit_set = self.clone();
let commit_bucket = bucket.to_owned();
let commit_object = object.to_owned();
let commit_upload_id = upload_id.to_owned();
let commit_upload_id_path = upload_id_path.clone();
let commit_version_suspended = opts.version_suspended;
let commit_is_versioned = opts.versioned || opts.version_suspended;
let commit_capacity_scope_token = opts.capacity_scope_token;
let commit_object_lock_guard = object_lock_guard.take();
let detach_commit_owner = commit_object_lock_guard.is_some() || upload_guard.is_some();
let commit = async move {
let _object_lock_guard = commit_object_lock_guard;
let complete_tail_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now);
// Crash-consistency injection: hard power loss after the upload is fully
// staged and locked but before the authoritative rename_data commit. No
// disk has moved the staged data, so a crash here must leave any prior
// committed version byte-for-byte intact (rustfs/backlog#864) and the
// upload fully retryable. Compiles to a no-op outside `#[cfg(test)]`.
if crash_inject::should_crash_at(CrashPoint::MultipartBeforeCommitRename, object) {
return Err(StorageError::Unexpected);
}
// The trailing `_` drops the rename_data old-size backfill
// (rustfs/backlog#1009): CompleteMultipartUpload keeps its pre-commit
// `get_object_info` lookup, so the backfill has no consumer here yet.
let (online_disks, convergence, op_old_dir, cleanup_disks, _) = Self::rename_data(
&shuffle_disks,
RUSTFS_META_MULTIPART_BUCKET,
&upload_id_path,
&parts_metadatas,
bucket,
object,
write_quorum,
)
.await?;
// Detach admission before any post-commit await: client cancellation
// must not couple durable convergence repair to cleanup work.
if convergence.needs_heal() {
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
false,
Some(HealChannelPriority::Normal),
Some(self.pool_index),
Some(self.set_index),
);
request.object_version_id = fi
.version_id
.or_else(|| opts.version_suspended.then(Uuid::nil))
.map(|version_id| version_id.to_string());
tokio::spawn(async move {
let _ = rustfs_common::heal_channel::send_heal_request(request).await;
});
}
// Crash-consistency injection: hard power loss after the authoritative
// rename_data commit succeeded but before the stale part.N.meta cleanup.
// The new version is durably committed and visible, so a crash here must
// leave the object readable as the new version; the un-reclaimed staging
// parts are swept by a retried completion or upload GC (rustfs/backlog#946).
// Compiles to a no-op outside `#[cfg(test)]`.
if crash_inject::should_crash_at(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object) {
return Err(StorageError::Unexpected);
}
// backlog#946: reclaim the stale per-part metadata (and any superfluous
// part.N data files no longer in the completed set) only *after* the
// authoritative rename_data commit above has succeeded. If rename_data
// fails write quorum and returns via `?`, the upload directory must keep
// its part.N.meta so a retried CompleteMultipartUpload can still read the
// parts; deleting them before the commit would strand the upload
// permanently. This mirrors the "clean up only after commit" pattern
// already used for the old data-dir GC and the upload-dir delete_all below.
self.cleanup_multipart_path(&parts).await;
if let Some(old_dir) = op_old_dir {
let committed_dir = fi.data_dir.unwrap_or_default().to_string();
// backlog#898: best-effort reclaim of the dereferenced old data dir.
// Returns a receipt (never `Err`); a failed GC must not turn an
// already-committed multipart completion into a 503.
let cleanup = self
.commit_rename_data_dir(&cleanup_disks, bucket, object, &old_dir.to_string(), &committed_dir, write_quorum)
.await;
self.report_old_data_dir_cleanup(bucket, object, &old_dir.to_string(), &cleanup)
.await;
}
if let Some(stage_start) = complete_tail_stage_start {
rustfs_io_metrics::record_put_object_stage_duration(
"multipart_complete_tail",
stage_start.elapsed().as_secs_f64() * 1000.0,
);
}
#[cfg(test)]
pause_multipart_commit(bucket, object, MultipartCommitPause::AfterRename).await;
let cleanup_store = self.clone();
let cleanup_upload_id_path = upload_id_path.clone();
let cleanup_bucket = bucket.to_owned();
let cleanup_object = object.to_owned();
let cleanup_upload_id = upload_id.to_owned();
let cleanup_handle = tokio::spawn(async move {
let _upload_guard = upload_guard;
let complete_tail_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now);
// Crash-consistency injection: hard power loss after the upload is fully
// staged and locked but before the authoritative rename_data commit. No
// disk has moved the staged data, so a crash here must leave any prior
// committed version byte-for-byte intact (rustfs/backlog#864) and the
// upload fully retryable. Compiles to a no-op outside `#[cfg(test)]`.
if crash_inject::should_crash_at(CrashPoint::MultipartBeforeCommitRename, &commit_object) {
return Err(StorageError::Unexpected);
}
// The trailing `_` drops the rename_data old-size backfill
// (rustfs/backlog#1009): CompleteMultipartUpload keeps its pre-commit
// `get_object_info` lookup, so the backfill has no consumer here yet.
let (online_disks, convergence, op_old_dir, cleanup_disks, _) = SetDisks::rename_data(
&shuffle_disks,
RUSTFS_META_MULTIPART_BUCKET,
&commit_upload_id_path,
&parts_metadatas,
&commit_bucket,
&commit_object,
write_quorum,
)
.await?;
// Detach admission before any post-commit await: client cancellation
// must not couple durable convergence repair to cleanup work.
if convergence.needs_heal() {
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
commit_bucket.clone(),
Some(commit_object.clone()),
false,
Some(HealChannelPriority::Normal),
Some(commit_set.pool_index),
Some(commit_set.set_index),
);
request.object_version_id = fi
.version_id
.or_else(|| commit_version_suspended.then(Uuid::nil))
.map(|version_id| version_id.to_string());
tokio::spawn(async move {
let _ = rustfs_common::heal_channel::send_heal_request(request).await;
});
}
// Crash-consistency injection: hard power loss after the authoritative
// rename_data commit succeeded but before the stale part.N.meta cleanup.
// The new version is durably committed and visible, so a crash here must
// leave the object readable as the new version; the un-reclaimed staging
// parts are swept by a retried completion or upload GC (rustfs/backlog#946).
// Compiles to a no-op outside `#[cfg(test)]`.
if crash_inject::should_crash_at(CrashPoint::MultipartAfterCommitBeforePartsCleanup, &commit_object) {
return Err(StorageError::Unexpected);
}
// backlog#946: reclaim the stale per-part metadata (and any superfluous
// part.N data files no longer in the completed set) only *after* the
// authoritative rename_data commit above has succeeded. If rename_data
// fails write quorum and returns via `?`, the upload directory must keep
// its part.N.meta so a retried CompleteMultipartUpload can still read the
// parts; deleting them before the commit would strand the upload
// permanently. This mirrors the "clean up only after commit" pattern
// already used for the old data-dir GC and the upload-dir delete_all below.
commit_set.cleanup_multipart_path(&parts).await;
if let Some(old_dir) = op_old_dir {
let committed_dir = fi.data_dir.unwrap_or_default().to_string();
// backlog#898: best-effort reclaim of the dereferenced old data dir.
// Returns a receipt (never `Err`); a failed GC must not turn an
// already-committed multipart completion into a 503.
let cleanup = commit_set
.commit_rename_data_dir(
&cleanup_disks,
&commit_bucket,
&commit_object,
&old_dir.to_string(),
&committed_dir,
write_quorum,
)
.await;
commit_set
.report_old_data_dir_cleanup(&commit_bucket, &commit_object, &old_dir.to_string(), &cleanup)
.await;
}
if let Some(stage_start) = complete_tail_stage_start {
rustfs_io_metrics::record_put_object_stage_duration(
"multipart_complete_tail",
stage_start.elapsed().as_secs_f64() * 1000.0,
);
}
#[cfg(test)]
pause_multipart_commit(&commit_bucket, &commit_object, MultipartCommitPause::AfterRename).await;
if let Err(err) = commit_set
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &commit_upload_id_path, write_quorum)
if let Err(err) = cleanup_store
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &cleanup_upload_id_path, write_quorum)
.await
{
warn!(
bucket = %commit_bucket,
object = %commit_object,
upload_id = %commit_upload_id,
bucket = %cleanup_bucket,
object = %cleanup_object,
upload_id = %cleanup_upload_id,
error = ?err,
"completed multipart upload staging cleanup did not reach write quorum"
);
}
for (i, op_disk) in online_disks.iter().enumerate() {
if let Some(disk) = op_disk
&& disk.is_online().await
{
fi = parts_metadatas[i].clone();
break;
}
}
commit_set.record_capacity_scope_if_needed(commit_capacity_scope_token, &online_disks);
fi.is_latest = true;
commit_set
.invalidate_get_object_metadata_cache(&commit_bucket, &commit_object)
.await;
drop(_object_lock_guard); // drop object lock guard to release the lock
drop(_upload_guard);
Ok(ObjectInfo::from_file_info(&fi, &commit_bucket, &commit_object, commit_is_versioned))
};
if detach_commit_owner {
tokio::spawn(commit)
.await
.map_err(|err| Error::other(format!("complete_multipart_upload commit task failed: {err}")))?
} else {
commit.await
});
if let Err(err) = cleanup_handle.await {
warn!(
bucket = %bucket,
object = %object,
upload_id = %upload_id,
error = ?err,
"completed multipart upload staging cleanup task failed"
);
}
drop(object_lock_guard); // drop object lock guard to release the lock
for (i, op_disk) in online_disks.iter().enumerate() {
if let Some(disk) = op_disk
&& disk.is_online().await
{
fi = parts_metadatas[i].clone();
break;
}
}
self.record_capacity_scope_if_needed(opts.capacity_scope_token, &online_disks);
fi.is_latest = true;
self.invalidate_get_object_metadata_cache(bucket, object).await;
Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended))
}
}
@@ -4900,74 +4883,6 @@ mod tests {
.await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn cancelled_complete_keeps_upload_lock_through_tail_cleanup() {
temp_env::async_with_vars(
[
(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true")),
(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("60")),
],
async {
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
let bucket = "multipart-cancelled-tail-lock-bucket";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
let (upload_id, parts) =
stage_upload_with_create_opts(&set_disks, bucket, object, &[0x53; 4096], &ObjectOptions::default()).await;
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path));
signaling.clear_observed();
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::AfterRename);
let complete_store = set_disks.clone();
let complete_upload_id = upload_id.clone();
let complete = tokio::spawn(async move {
complete_store
.complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default())
.await
});
barrier.wait_until_paused().await;
let abort_store = set_disks.clone();
let abort_upload_id = upload_id.clone();
let abort = tokio::spawn(async move {
abort_store
.abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default())
.await
});
signaling.wait_for_attempts(2).await;
tokio::task::yield_now().await;
assert!(!abort.is_finished(), "abort must wait while completion tail owns the upload lock");
complete.abort();
assert!(
complete
.await
.expect_err("the completion request should be cancellable while the tail is paused")
.is_cancelled()
);
tokio::task::yield_now().await;
assert!(!abort.is_finished(), "cancelling the completion waiter must not release the upload lock");
barrier.release();
let abort_err = abort
.await
.expect("abort task should not panic")
.expect_err("the committed upload should no longer exist when abort acquires the lock");
assert!(
matches!(abort_err, StorageError::InvalidUploadID(..)),
"abort should return InvalidUploadID after the detached completion tail, got {abort_err:?}"
);
},
)
.await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn complete_validates_parts_after_an_inflight_upload_part_commit() {
+154 -516
View File
@@ -188,74 +188,6 @@ async fn get_object_reader_with_context(
GetObjectReader::new_with_resolver(reader, range, object_info, opts, headers, ctx.object_encryption_resolver()).await
}
fn data_read_metadata_early_stop_request_shape_allowed(range: &Option<HTTPRangeSpec>, opts: &ObjectOptions) -> bool {
range.is_none()
&& opts.part_number.is_none()
&& opts.version_id.is_none()
&& !opts.incl_free_versions
&& !opts.skip_free_version
&& !opts.raw_data_movement_read
&& !opts.data_movement
&& !crate::object_api::restore_request_active(opts)
}
#[cfg(test)]
mod data_read_metadata_early_stop_request_shape_tests {
use super::*;
#[test]
fn data_read_metadata_early_stop_only_allows_whole_latest_plain_get_shape() {
assert!(data_read_metadata_early_stop_request_shape_allowed(&None, &ObjectOptions::default()));
let range = Some(HTTPRangeSpec {
is_suffix_length: false,
start: 0,
end: 0,
});
assert!(!data_read_metadata_early_stop_request_shape_allowed(&range, &ObjectOptions::default()));
let part_opts = ObjectOptions {
part_number: Some(1),
..Default::default()
};
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &part_opts));
let version_opts = ObjectOptions {
version_id: Some(Uuid::new_v4().to_string()),
..Default::default()
};
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &version_opts));
let incl_free_opts = ObjectOptions {
incl_free_versions: true,
..Default::default()
};
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &incl_free_opts));
let skip_free_opts = ObjectOptions {
skip_free_version: true,
..Default::default()
};
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &skip_free_opts));
let data_movement_opts = ObjectOptions {
data_movement: true,
..Default::default()
};
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &data_movement_opts));
let raw_data_movement_opts = ObjectOptions {
raw_data_movement_read: true,
..Default::default()
};
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &raw_data_movement_opts));
let mut restore_opts = ObjectOptions::default();
restore_opts.transition.restore_request.days = Some(1);
assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &restore_opts));
}
}
/// Length of the full plaintext body when — and only when — this read's output
/// is exactly the object's complete plaintext, so the app-layer body cache may
/// serve it in place of the erasure read.
@@ -499,16 +431,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
let (snapshot, prepared_object_info) = if let Some(prepared) = take_prepared_get_object_metadata() {
(prepared.snapshot, prepared.object_info)
} else {
match self
.get_object_fileinfo(
bucket,
object,
opts,
true,
data_read_metadata_early_stop_request_shape_allowed(&range, opts),
)
.await
{
match self.get_object_fileinfo(bucket, object, opts, true, true).await {
Ok(snapshot) => (snapshot, None),
Err(err) => {
rustfs_io_metrics::record_get_object_metadata_phase_duration(metadata_stage_start.elapsed().as_secs_f64());
@@ -660,7 +583,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
}
}
let erasure = self.erasure_cache.get_for_file_info(fi)?;
let erasure = erasure_from_file_info(fi, fi.uses_legacy_checksum)?;
let read_length = erasure.shard_file_offset(0, object_size, object_size);
let total_shards = data_shards + fi.erasure.parity_blocks;
let (_disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, files, fi);
@@ -829,7 +752,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
if let Some(body) = Self::try_get_object_direct_data_shards_with_fileinfo(
bucket,
object,
Arc::clone(&self.erasure_cache),
fi,
files,
disks,
@@ -865,7 +787,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
Self::get_object_with_fileinfo(
bucket,
object,
Arc::clone(&self.erasure_cache),
0,
object_info.size,
&mut output,
@@ -909,7 +830,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
match Self::get_object_decode_reader_with_fileinfo(
bucket,
object,
Arc::clone(&self.erasure_cache),
fi,
files,
disks,
@@ -974,7 +894,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
let set_index = self.set_index;
let pool_index = self.pool_index;
let skip_verify = opts.skip_verify_bitrot;
let erasure_cache = Arc::clone(&self.erasure_cache);
let (fi, files, disks) = snapshot.into_owned();
tokio::spawn(async move {
let _guard = read_lock_guard;
@@ -986,7 +905,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
if let Err(e) = Self::get_object_with_fileinfo(
&bucket,
&object,
erasure_cache,
offset,
length,
&mut writer,
@@ -1188,7 +1106,6 @@ impl SetDisks {
let tmp_object = format!("{}/{}/part.1", tmp_dir, fi.data_dir.unwrap());
let mut tmp_cleanup_owned = false;
let result: Result<(ObjectInfo, Option<OldCurrentSize>)> = async {
let erasure = Arc::new(erasure_from_file_info(&fi, false)?);
@@ -1680,236 +1597,169 @@ impl SetDisks {
});
}
let commit_set = self.clone();
let commit_bucket = bucket.to_owned();
let commit_object = object.to_owned();
let commit_tmp_dir = tmp_dir.clone();
let commit_object_lock_guard = object_lock_guard.take();
let commit_bucket_lifecycle_guard = bucket_lifecycle_guard.take();
let detach_commit_owner = commit_object_lock_guard.is_some() || commit_bucket_lifecycle_guard.is_some();
let commit_write_path_label = write_path.metric_label();
let commit_is_versioned = opts.versioned || opts.version_suspended;
let commit_capacity_scope_token = opts.capacity_scope_token;
let commit_replication_state = replication_state_to_filemeta(&opts.put_replication_state());
tmp_cleanup_owned = true;
let rename_stage_start = Instant::now();
let (online_disks, convergence, op_old_dir, cleanup_disks, old_current_size) = Self::rename_data(
&shuffle_disks,
RUSTFS_META_TMP_BUCKET,
tmp_dir.as_str(),
&parts_metadatas,
bucket,
object,
write_quorum,
)
.await?;
// Do this before any post-commit await so request cancellation cannot
// bypass best-effort admission. A process crash before admission
// remains subject to the existing scanner reconciliation path.
if convergence.needs_heal() {
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
false,
Some(HealChannelPriority::Normal),
Some(self.pool_index),
Some(self.set_index),
);
request.object_version_id = committed_version_id.map(|version_id| version_id.to_string());
tokio::spawn(async move {
let _ = rustfs_common::heal_channel::send_heal_request(request).await;
});
}
let commit = async move {
let _object_lock_guard = commit_object_lock_guard;
let _bucket_lifecycle_guard = commit_bucket_lifecycle_guard;
let rename_stage_start = Instant::now();
let rename_result = SetDisks::rename_data(
&shuffle_disks,
RUSTFS_META_TMP_BUCKET,
commit_tmp_dir.as_str(),
&parts_metadatas,
&commit_bucket,
&commit_object,
let rename_stage_elapsed = rename_stage_start.elapsed();
let rename_stage_ms = rename_stage_elapsed.as_millis() as u64;
self.invalidate_get_object_metadata_cache(bucket, object).await;
// `rename_data` has completed the authoritative quorum commit. The
// exact old-data-dir reclamation below is best-effort space cleanup;
// it must not serialize the next operation on this object.
drop(object_lock_guard);
rustfs_io_metrics::record_put_object_stage_duration("set_disk_rename", duration_millis_f64(rename_stage_elapsed));
if (rename_stage_ms as u128) >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS {
warn!(
event = EVENT_SET_DISK_COMMIT_TAIL_SLOW,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
stage = "rename_data",
bucket = %bucket,
object = %object,
tmp_dir = %tmp_dir,
duration_ms = { rename_stage_ms },
write_quorum,
)
.await;
let (online_disks, convergence, op_old_dir, cleanup_disks, old_current_size) = match rename_result {
Ok(commit) => commit,
Err(err) => {
if let Err(cleanup_err) = commit_set.delete_all(RUSTFS_META_TMP_BUCKET, &commit_tmp_dir).await {
warn!(tmp_dir = %commit_tmp_dir, error = ?cleanup_err, "failed to cleanup put_object temporary data");
} else if issue3031_diag_enabled() {
warn!(
target: "rustfs_ecstore::set_disk",
bucket = %commit_bucket,
object = %commit_object,
tmp_dir = %commit_tmp_dir,
"issue3031_put_object_tmp_cleanup_done"
);
}
return Err(err.into());
}
};
// Do this before any post-commit await so request cancellation cannot
// bypass best-effort admission. A process crash before admission
// remains subject to the existing scanner reconciliation path.
if convergence.needs_heal() {
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
commit_bucket.clone(),
Some(commit_object.clone()),
false,
Some(HealChannelPriority::Normal),
Some(commit_set.pool_index),
Some(commit_set.set_index),
);
request.object_version_id = committed_version_id.map(|version_id| version_id.to_string());
tokio::spawn(async move {
let _ = rustfs_common::heal_channel::send_heal_request(request).await;
});
}
state = "slow",
"SetDisk commit tail stage is slow"
);
}
let rename_stage_elapsed = rename_stage_start.elapsed();
let rename_stage_ms = rename_stage_elapsed.as_millis() as u64;
commit_set
.invalidate_get_object_metadata_cache(&commit_bucket, &commit_object)
let mut cleanup_stage_ms: Option<u64> = None;
if let Some(old_dir) = op_old_dir {
let committed_dir = committed_data_dir.unwrap_or_default().to_string();
let cleanup_stage_start = Instant::now();
// backlog#898: reclaiming the dereferenced old data dir is
// best-effort and returns a receipt (never `Err`). A failed GC
// here must not negate an already-committed, durable write, so we
// deliberately do NOT `?`-propagate it into a 503. On residue the
// report path emits the leak metric and enqueues a heal.
let cleanup = self
.commit_rename_data_dir(&cleanup_disks, bucket, object, &old_dir.to_string(), &committed_dir, write_quorum)
.await;
// `rename_data` has completed the authoritative quorum commit. The
// exact old-data-dir reclamation below is best-effort space cleanup;
// it must not serialize the next operation on this object.
drop(_object_lock_guard);
drop(_bucket_lifecycle_guard);
rustfs_io_metrics::record_put_object_stage_duration("set_disk_rename", duration_millis_f64(rename_stage_elapsed));
if (rename_stage_ms as u128) >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS {
let cleanup_elapsed = cleanup_stage_start.elapsed();
let cleanup_ms = cleanup_elapsed.as_millis() as u64;
cleanup_stage_ms = Some(cleanup_ms);
rustfs_io_metrics::record_put_object_stage_duration(
"set_disk_old_data_cleanup",
duration_millis_f64(cleanup_elapsed),
);
self.report_old_data_dir_cleanup(bucket, object, &old_dir.to_string(), &cleanup)
.await;
if (cleanup_ms as u128) >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS {
warn!(
event = EVENT_SET_DISK_COMMIT_TAIL_SLOW,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
stage = "rename_data",
bucket = %commit_bucket,
object = %commit_object,
tmp_dir = %commit_tmp_dir,
duration_ms = { rename_stage_ms },
stage = "commit_rename_data_dir",
bucket = %bucket,
object = %object,
tmp_dir = %tmp_dir,
old_dir = %old_dir,
duration_ms = cleanup_ms,
write_quorum,
state = "slow",
"SetDisk commit tail stage is slow"
);
}
let mut cleanup_stage_ms: Option<u64> = None;
if let Some(old_dir) = op_old_dir {
let committed_dir = committed_data_dir.unwrap_or_default().to_string();
let cleanup_stage_start = Instant::now();
// backlog#898: reclaiming the dereferenced old data dir is
// best-effort and returns a receipt (never `Err`). A failed GC
// here must not negate an already-committed, durable write, so we
// deliberately do NOT `?`-propagate it into a 503. On residue the
// report path emits the leak metric and enqueues a heal.
let cleanup = commit_set
.commit_rename_data_dir(
&cleanup_disks,
&commit_bucket,
&commit_object,
&old_dir.to_string(),
&committed_dir,
write_quorum,
)
.await;
let cleanup_elapsed = cleanup_stage_start.elapsed();
let cleanup_ms = cleanup_elapsed.as_millis() as u64;
cleanup_stage_ms = Some(cleanup_ms);
rustfs_io_metrics::record_put_object_stage_duration(
"set_disk_old_data_cleanup",
duration_millis_f64(cleanup_elapsed),
);
commit_set
.report_old_data_dir_cleanup(&commit_bucket, &commit_object, &old_dir.to_string(), &cleanup)
.await;
if (cleanup_ms as u128) >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS {
warn!(
event = EVENT_SET_DISK_COMMIT_TAIL_SLOW,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
stage = "commit_rename_data_dir",
bucket = %commit_bucket,
object = %commit_object,
tmp_dir = %commit_tmp_dir,
old_dir = %old_dir,
duration_ms = cleanup_ms,
write_quorum,
state = "slow",
"SetDisk commit tail stage is slow"
);
}
}
let committed_metadata_slot = committed_response_metadata_slot(&online_disks, response_metadata_slot);
let mut fi = std::mem::take(&mut parts_metadatas[committed_metadata_slot]);
if is_compressed {
record_compression_total_memory(actual_size as u64, w_size as u64).await;
}
commit_set.record_capacity_scope_if_needed(commit_capacity_scope_token, &online_disks);
fi.replication_state_internal = Some(commit_replication_state);
fi.is_latest = true;
if issue3031_diag_enabled() {
let online_success_count = online_disks.iter().filter(|disk| disk.is_some()).count();
warn!(
target: "rustfs_ecstore::set_disk",
bucket = %commit_bucket,
object = %commit_object,
tmp_dir = %commit_tmp_dir,
data_dir = ?fi.data_dir,
write_quorum,
online_success_count,
op_old_dir = ?op_old_dir,
"issue3031_put_object_commit_succeeded"
);
}
let total_commit_tail_ms = rename_stage_start.elapsed().as_millis();
if total_commit_tail_ms >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS {
warn!(
event = EVENT_SET_DISK_COMMIT_TAIL_SLOW,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
stage = "put_object_commit_tail",
bucket = %commit_bucket,
object = %commit_object,
tmp_dir = %commit_tmp_dir,
duration_ms = total_commit_tail_ms as u64,
write_quorum,
state = "slow",
"SetDisk commit tail is slow"
);
}
if issue3031_diag_enabled() {
warn!(
event = EVENT_SET_DISK_PUT_OBJECT_STAGE_SUMMARY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket = %commit_bucket,
object = %commit_object,
write_quorum,
write_path = commit_write_path_label,
writer_setup_ms,
encode_ms,
rename_ms = rename_stage_ms,
cleanup_ms = cleanup_stage_ms.unwrap_or_default(),
cleanup_present = cleanup_stage_ms.is_some(),
commit_tail_ms = total_commit_tail_ms as u64,
result = "success",
"SetDisk put_object stage summary"
);
}
let cleanup_set = commit_set.clone();
let cleanup_tmp_dir = commit_tmp_dir.clone();
tokio::spawn(async move {
if let Err(err) = cleanup_set.delete_all(RUSTFS_META_TMP_BUCKET, &cleanup_tmp_dir).await {
warn!(tmp_dir = %cleanup_tmp_dir, error = ?err, "failed to cleanup put_object temporary data");
} else if issue3031_diag_enabled() {
warn!(
target: "rustfs_ecstore::set_disk",
tmp_dir = %cleanup_tmp_dir,
"issue3031_put_object_tmp_cleanup_done"
);
}
});
Ok((
ObjectInfo::from_file_info(&fi, &commit_bucket, &commit_object, commit_is_versioned),
old_current_size,
))
};
if detach_commit_owner {
tokio::spawn(commit)
.await
.map_err(|err| Error::other(format!("put_object commit task failed: {err}")))?
} else {
commit.await
}
let committed_metadata_slot = committed_response_metadata_slot(&online_disks, response_metadata_slot);
let mut fi = std::mem::take(&mut parts_metadatas[committed_metadata_slot]);
if is_compressed {
record_compression_total_memory(actual_size as u64, w_size as u64).await;
}
self.record_capacity_scope_if_needed(opts.capacity_scope_token, &online_disks);
fi.replication_state_internal = Some(replication_state_to_filemeta(&opts.put_replication_state()));
fi.is_latest = true;
if issue3031_diag_enabled() {
let online_success_count = online_disks.iter().filter(|disk| disk.is_some()).count();
warn!(
target: "rustfs_ecstore::set_disk",
bucket = %bucket,
object = %object,
tmp_dir = %tmp_dir,
data_dir = ?fi.data_dir,
write_quorum,
online_success_count,
op_old_dir = ?op_old_dir,
"issue3031_put_object_commit_succeeded"
);
}
let total_commit_tail_ms = rename_stage_start.elapsed().as_millis();
if total_commit_tail_ms >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS {
warn!(
event = EVENT_SET_DISK_COMMIT_TAIL_SLOW,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
stage = "put_object_commit_tail",
bucket = %bucket,
object = %object,
tmp_dir = %tmp_dir,
duration_ms = total_commit_tail_ms as u64,
write_quorum,
state = "slow",
"SetDisk commit tail is slow"
);
}
if issue3031_diag_enabled() {
warn!(
event = EVENT_SET_DISK_PUT_OBJECT_STAGE_SUMMARY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket = %bucket,
object = %object,
write_quorum,
write_path = write_path.metric_label(),
writer_setup_ms,
encode_ms,
rename_ms = rename_stage_ms,
cleanup_ms = cleanup_stage_ms.unwrap_or_default(),
cleanup_present = cleanup_stage_ms.is_some(),
commit_tail_ms = total_commit_tail_ms as u64,
result = "success",
"SetDisk put_object stage summary"
);
}
Ok((
ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended),
old_current_size,
))
}
.await;
@@ -1945,8 +1795,7 @@ impl SetDisks {
);
}
if tmp_cleanup_owned && result.is_ok() {
} else if result.is_ok() {
if result.is_ok() {
// Success path: `rename_data` has already moved the data dir out of
// the tmp workspace and removed the (empty) tmp dir where it could,
// so this delete_all is a speculative safety net that normally hits
@@ -5013,13 +4862,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let pool_index = self.pool_index;
let skip_verify = opts.skip_verify_bitrot;
let metrics_size_bucket = rustfs_io_metrics::get_object_size_bucket(cloned_fi.size);
let erasure_cache = Arc::clone(&self.erasure_cache);
let producer = async move {
let mut writer = TransitionUploadWriter::new(pw);
Self::get_object_with_fileinfo(
&cloned_bucket,
&cloned_object,
erasure_cache,
0,
cloned_fi.size,
&mut writer,
@@ -5993,37 +5840,6 @@ mod inline_put_commit_path_tests {
assert_eq!(restored, payload);
}
#[tokio::test]
async fn repeated_gets_reuse_the_set_erasure_shell() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "get-erasure-shell-cache";
let object = "object.bin";
let payload = vec![0x4d; 1024 * 1024];
make_bucket(&disk_stores, bucket).await;
let mut reader = PutObjReader::from_vec(payload.clone());
set_disks
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("non-inline object should commit");
assert!(set_disks.erasure_cache.entries.read().is_empty());
for _ in 0..2 {
let mut object_reader = set_disks
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("cached-shell GET should succeed");
let mut restored = Vec::new();
object_reader
.stream
.read_to_end(&mut restored)
.await
.expect("cached-shell GET should stream");
assert_eq!(restored, payload);
assert_eq!(set_disks.erasure_cache.entries.read().len(), 1);
}
}
#[tokio::test]
async fn inline_put_direct_commit_accepts_exact_quorum_and_rejects_quorum_minus_one() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
@@ -6217,10 +6033,7 @@ mod inline_put_commit_path_tests {
mod get_object_downstream_close_accounting_tests {
use super::hermetic_set_disks_support::hermetic_set_disks;
use super::*;
use crate::diagnostics::get::{
GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, GET_OBJECT_PATH_INTERNAL_META, GET_STAGE_DECODE, GET_STAGE_EMIT,
GetObjectFailureReason,
};
use crate::diagnostics::get::{GET_OBJECT_PATH_INTERNAL_META, GET_STAGE_DECODE, GET_STAGE_EMIT, GetObjectFailureReason};
use crate::disk::RUSTFS_META_BUCKET;
use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions};
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
@@ -6339,22 +6152,7 @@ mod get_object_downstream_close_accounting_tests {
let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled();
rustfs_io_metrics::set_get_stage_metrics_enabled(true);
let (
internal_missing,
legacy_unknown,
internal_fanout,
legacy_fanout,
internal_scheduled,
legacy_scheduled,
internal_completed,
legacy_completed,
internal_cancelled,
legacy_cancelled,
internal_unsafe_miss,
legacy_unsafe_miss,
internal_saved,
legacy_saved,
) = metrics::with_local_recorder(&recorder, || {
let (internal_missing, legacy_unknown, internal_fanout, legacy_fanout) = metrics::with_local_recorder(&recorder, || {
runtime.block_on(async {
let (_temp_dirs, _disk_stores, set_disks) = hermetic_set_disks(4).await;
let options = ObjectOptions {
@@ -6398,54 +6196,6 @@ mod get_object_downstream_close_accounting_tests {
"rustfs_io_get_object_metadata_fanout_error_responses",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_scheduled",
&[("path", GET_OBJECT_PATH_INTERNAL_META)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_scheduled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_completed",
&[("path", GET_OBJECT_PATH_INTERNAL_META)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_completed",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_cancelled",
&[("path", GET_OBJECT_PATH_INTERNAL_META)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_fanout_cancelled",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)],
),
recorder.counter_value(
"rustfs_io_get_object_metadata_early_stop_total",
&[
("path", GET_OBJECT_PATH_INTERNAL_META),
("decision", "miss"),
("reason", GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST),
],
),
recorder.counter_value(
"rustfs_io_get_object_metadata_early_stop_total",
&[
("path", GET_OBJECT_PATH_LEGACY_DUPLEX),
("decision", "miss"),
("reason", GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST),
],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_early_stop_saved_responses",
&[("path", GET_OBJECT_PATH_INTERNAL_META)],
),
recorder.histogram_values(
"rustfs_io_get_object_metadata_early_stop_saved_responses",
&[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)],
),
)
})
});
@@ -6458,50 +6208,6 @@ mod get_object_downstream_close_accounting_tests {
);
assert_eq!(internal_fanout, vec![4.0], "internal metadata fanout must retain its path label");
assert!(legacy_fanout.is_empty(), "internal metadata fanout must not leak into legacy_duplex");
assert_eq!(
internal_scheduled,
vec![4.0],
"internal metadata lifecycle scheduled count must retain its path label"
);
assert!(
legacy_scheduled.is_empty(),
"internal metadata lifecycle scheduled count must not leak into legacy_duplex"
);
assert_eq!(
internal_completed,
vec![4.0],
"internal metadata lifecycle completed count must retain its path label"
);
assert!(
legacy_completed.is_empty(),
"internal metadata lifecycle completed count must not leak into legacy_duplex"
);
assert_eq!(
internal_cancelled,
vec![0.0],
"internal metadata full-wait lifecycle must record zero cancellations"
);
assert!(
legacy_cancelled.is_empty(),
"internal metadata lifecycle cancelled count must not leak into legacy_duplex"
);
assert_eq!(
internal_unsafe_miss, 1,
"internal metadata unsafe early-stop miss must retain its path label"
);
assert_eq!(
legacy_unsafe_miss, 0,
"internal metadata unsafe early-stop miss must not leak into legacy_duplex"
);
assert_eq!(
internal_saved,
vec![0.0],
"internal metadata unsafe miss must record zero saved responses on internal_meta"
);
assert!(
legacy_saved.is_empty(),
"internal metadata unsafe miss saved responses must not leak into legacy_duplex"
);
}
}
@@ -10056,74 +9762,6 @@ mod put_object_tmp_cleanup_tests {
assert_eq!(body, vec![b'2'; TEST_OBJECT_SIZE]);
}
#[tokio::test]
async fn cancelled_rename_keeps_namespace_lock_until_publication() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "put-commit-lock-cancelled-rename";
let object = "commit-lock-cancelled-rename-object";
for disk in &disk_stores {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
let rename_tasks = rename_fanout_barrier::observe_tasks(object);
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
let first_store = Arc::clone(&set_disks);
let first = tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(vec![b'1'; TEST_OBJECT_SIZE]);
first_store
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
});
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
.await
.expect("first PUT should pause during the authoritative rename");
let second_namespace_barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::BeforeNamespace);
let second_store = Arc::clone(&set_disks);
let second = tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(vec![b'2'; TEST_OBJECT_SIZE]);
second_store
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
});
second_namespace_barrier.release_and_wait_until_namespace_pending().await;
first.abort();
assert!(
first
.await
.expect_err("the first request should be cancelled while rename is parked")
.is_cancelled()
);
tokio::task::yield_now().await;
assert!(
!second.is_finished(),
"the second writer must remain blocked by the cancelled commit owner"
);
rename_barrier.release();
drop(rename_barrier);
tokio::time::timeout(Duration::from_secs(30), async {
while rename_tasks.running() != 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("the cancelled owner's rename fanout should drain");
second
.await
.expect("second overwrite task should join")
.expect("second overwrite should commit after the cancelled owner reaches publication");
let mut reader = set_disks
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("the latest overwrite should be readable");
let mut body = Vec::new();
reader.stream.read_to_end(&mut body).await.expect("latest body should drain");
assert_eq!(body, vec![b'2'; TEST_OBJECT_SIZE]);
}
#[tokio::test]
async fn put_object_no_lock_aborts_after_outer_namespace_lock_loss() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
+23 -22
View File
@@ -482,7 +482,6 @@ impl SetDisks {
pub(super) async fn try_get_object_direct_data_shards_with_fileinfo(
bucket: &str,
object: &str,
erasure_cache: Arc<ErasureCache>,
fi: &FileInfo,
files: &[FileInfo],
disks: &[Option<DiskStore>],
@@ -503,7 +502,13 @@ impl SetDisks {
return Ok(None);
}
let erasure = erasure_cache.get_for_file_info(fi)?;
let erasure = coding::Erasure::try_new_with_options(
fi.erasure.data_blocks,
fi.erasure.parity_blocks,
fi.erasure.block_size,
fi.uses_legacy_checksum,
)
.map_err(Error::from)?;
let checksum_info = fi.erasure.get_checksum_info(part.number);
let checksum_algo = if fi.uses_legacy_checksum && checksum_info.algorithm == HashAlgorithm::HighwayHash256S {
@@ -631,7 +636,6 @@ impl SetDisks {
// &self,
bucket: &str,
object: &str,
erasure_cache: Arc<ErasureCache>,
offset: usize,
length: i64,
writer: &mut W,
@@ -726,7 +730,13 @@ impl SetDisks {
object, offset, length, end_offset, part_index, last_part_index, last_part_relative_offset, "Multipart read bounds"
);
let erasure = erasure_cache.get_for_file_info(&fi)?;
let erasure = coding::Erasure::try_new_with_options(
fi.erasure.data_blocks,
fi.erasure.parity_blocks,
fi.erasure.block_size,
fi.uses_legacy_checksum,
)
.map_err(Error::from)?;
let part_indices: Vec<usize> = (part_index..=last_part_index).collect();
debug!(bucket, object, ?part_indices, "Multipart part indices to stream");
@@ -1160,7 +1170,6 @@ impl SetDisks {
pub(super) async fn get_object_decode_reader_with_fileinfo(
bucket: &str,
object: &str,
erasure_cache: Arc<ErasureCache>,
fi: &FileInfo,
files: &[FileInfo],
disks: &[Option<DiskStore>],
@@ -1171,7 +1180,14 @@ impl SetDisks {
metrics_size_bucket: &'static str,
prefer_data_blocks_first_reader_setup: bool,
) -> Result<GetCodecStreamingReaderBuildOutcome> {
let erasure = erasure_cache.get_for_file_info(fi)?;
let erasure = coding::Erasure::try_new_with_options(
fi.erasure.data_blocks,
fi.erasure.parity_blocks,
fi.erasure.block_size,
fi.uses_legacy_checksum,
)
.map_err(Error::from)?;
let (disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, files, fi);
if fi.parts.len() == 1 {
@@ -1558,7 +1574,7 @@ struct LazyCodecPartContext {
fi: FileInfo,
files: Vec<FileInfo>,
disks: Vec<Option<DiskStore>>,
erasure: Arc<coding::Erasure>,
erasure: coding::Erasure,
skip_verify_bitrot: bool,
metrics_object_class: &'static str,
metrics_size_bucket: &'static str,
@@ -2042,7 +2058,6 @@ mod metadata_cache_tests {
let err = SetDisks::get_object_with_fileinfo(
"bucket",
"object",
Arc::new(ErasureCache::new()),
0,
1,
&mut output,
@@ -2073,7 +2088,6 @@ mod metadata_cache_tests {
let err = SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
2,
1,
&mut output,
@@ -2097,7 +2111,6 @@ mod metadata_cache_tests {
let err = SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
usize::MAX,
1,
&mut output,
@@ -2119,7 +2132,6 @@ mod metadata_cache_tests {
let err = SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
1,
1,
&mut output,
@@ -2143,7 +2155,6 @@ mod metadata_cache_tests {
let err = SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
0,
1,
&mut output,
@@ -2181,7 +2192,6 @@ mod metadata_cache_tests {
SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
0,
0,
&mut output,
@@ -2214,7 +2224,6 @@ mod metadata_cache_tests {
let err = SetDisks::get_object_with_fileinfo(
bucket,
object,
Arc::new(ErasureCache::new()),
0,
1,
&mut output,
@@ -4119,7 +4128,6 @@ mod tests {
let result = SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&fi,
&[],
&[],
@@ -4142,7 +4150,6 @@ mod tests {
let invalid_size = SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&single_part,
&[],
&[],
@@ -4163,7 +4170,6 @@ mod tests {
SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&multipart,
&[],
&[],
@@ -4188,7 +4194,6 @@ mod tests {
SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&multipart,
&[],
&[],
@@ -4217,7 +4222,6 @@ mod tests {
SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&multipart,
&[],
&[],
@@ -4271,7 +4275,6 @@ mod tests {
SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&fi,
&files,
&disks,
@@ -4325,7 +4328,6 @@ mod tests {
SetDisks::get_object_decode_reader_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
&fi,
&files,
&disks,
@@ -4370,7 +4372,6 @@ mod tests {
SetDisks::get_object_with_fileinfo(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
Arc::new(ErasureCache::new()),
0,
part_data.len() as i64,
&mut output,
-2
View File
@@ -14,8 +14,6 @@
//! test endpoint index settings
#![recursion_limit = "256"]
use std::net::SocketAddr;
use tempfile::TempDir;
use tokio_util::sync::CancellationToken;
@@ -22,8 +22,6 @@
//! bucket-metadata-sys OnceCell) — under `cargo nextest` each test runs
//! in its own process so the OnceCell never collides.
#![recursion_limit = "256"]
use http::HeaderMap;
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
use rustfs_heal::heal::{
@@ -21,8 +21,6 @@
//! These drive the REAL `ECStoreHealStorage` + `ECStore` against real disks.
//! Every test is `#[serial]`; under `cargo nextest` each runs in its own process.
#![recursion_limit = "256"]
use http::HeaderMap;
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
use rustfs_heal::heal::storage::{
@@ -12,8 +12,6 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#![recursion_limit = "256"]
use http::HeaderMap;
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
use rustfs_heal::heal::{
-44
View File
@@ -812,17 +812,6 @@ pub fn record_get_object_metadata_fanout_shape(path: &'static str, total: usize,
.record(metadata_fanout_count_to_f64(non_valid));
}
/// Record task lifecycle shape for one GetObject metadata fanout.
#[inline(always)]
pub fn record_get_object_metadata_fanout_lifecycle(path: &'static str, scheduled: usize, completed: usize, cancelled: usize) {
if !get_stage_metrics_enabled() {
return;
}
histogram!("rustfs_io_get_object_metadata_fanout_scheduled", "path" => path).record(metadata_fanout_count_to_f64(scheduled));
histogram!("rustfs_io_get_object_metadata_fanout_completed", "path" => path).record(metadata_fanout_count_to_f64(completed));
histogram!("rustfs_io_get_object_metadata_fanout_cancelled", "path" => path).record(metadata_fanout_count_to_f64(cancelled));
}
/// Record a guarded metadata early-stop hit for GetObject.
#[inline(always)]
pub fn record_get_object_metadata_early_stop_hit(path: &'static str, reason: &'static str) {
@@ -2709,7 +2698,6 @@ mod tests {
record_get_object_quorum_reached_latency("legacy_duplex", 0.002);
record_get_object_metadata_response("legacy_duplex", "valid");
record_get_object_metadata_fanout_shape("legacy_duplex", 4, 3, 1, 1);
record_get_object_metadata_fanout_lifecycle("legacy_duplex", 4, 3, 1);
record_get_object_metadata_early_stop_hit("legacy_duplex", "valid_quorum");
record_get_object_metadata_early_stop_miss("legacy_duplex", "insufficient_quorum");
record_get_object_metadata_early_stop_saved_responses("legacy_duplex", 1);
@@ -2780,38 +2768,6 @@ mod tests {
assert!(remote_scheduled >= remote_avoid_potential);
}
#[test]
fn metadata_fanout_lifecycle_records_named_histograms() {
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
set_get_stage_metrics_enabled(true);
record_get_object_metadata_fanout_lifecycle("legacy_duplex", 4, 3, 1);
set_get_stage_metrics_enabled(false);
});
let metrics = snapshotter.snapshot().into_vec();
for (name, expected) in [
("rustfs_io_get_object_metadata_fanout_scheduled", 4.0),
("rustfs_io_get_object_metadata_fanout_completed", 3.0),
("rustfs_io_get_object_metadata_fanout_cancelled", 1.0),
] {
let value = metrics.iter().find_map(|(composite, _, _, value)| {
let has_path = composite
.key()
.labels()
.any(|label| label.key() == "path" && label.value() == "legacy_duplex");
(composite.kind() == MetricKind::Histogram && composite.key().name() == name && has_path).then_some(value)
});
assert!(
matches!(value, Some(DebugValue::Histogram(values)) if values.len() == 1 && values[0].0 == expected),
"{name} must record the exact fanout lifecycle sample"
);
}
}
#[test]
fn test_record_get_object_fill_metrics() {
record_get_object_fill_queued("codec_streaming", "single_inflight", 1);
@@ -23,7 +23,6 @@
//! two are tested together because a reload is the only way to tell a real
//! merge from one that happened to look right in the cache.
#![recursion_limit = "256"]
#![cfg(feature = "swift")]
use std::collections::HashMap;
-1
View File
@@ -12,7 +12,6 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#![recursion_limit = "256"]
#![cfg_attr(docsrs, feature(doc_auto_cfg))]
#![warn(
// missing_docs,
@@ -12,8 +12,6 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#![recursion_limit = "256"]
use futures::FutureExt;
use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT;
use rustfs_scanner::scanner_folder::ScannerItem;
+92 -299
View File
@@ -12,6 +12,7 @@ use datafusion::{
use std::sync::Arc;
use crate::table_catalog::test_support::{
TestCatalogObjectBackend as TestTableCatalogObjectBackend, TestCatalogObjectRecord,
manifest_avro_bytes as test_manifest_avro_bytes,
manifest_avro_bytes_with_nullable_sequences as test_manifest_avro_bytes_with_nullable_sequences,
manifest_list_avro_bytes as test_manifest_list_avro_bytes, manifest_list_avro_entries as test_manifest_list_avro_entries,
@@ -2017,7 +2018,7 @@ fn format_upgrade_assigns_v1_snapshot_sequences_and_rejects_v3() {
#[tokio::test]
async fn create_table_response_writes_initial_metadata_for_standard_request() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -2094,7 +2095,7 @@ async fn create_table_holds_bucket_fence_from_metadata_write_through_registratio
let barrier = Arc::new(tokio::sync::Barrier::new(2));
let metadata_backend = TestTableCatalogObjectBackend {
put_object_barrier: Some(Arc::clone(&barrier)),
..Default::default()
..TestTableCatalogObjectBackend::content_addressed()
};
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(store.as_ref(), "warehouse", true)
@@ -2128,7 +2129,7 @@ async fn create_table_holds_bucket_fence_from_metadata_write_through_registratio
create_table_response(create_store.as_ref(), &create_backend, "warehouse", &create_namespace, request, true).await
});
tokio::time::timeout(StdDuration::from_secs(2), async {
while metadata_backend.objects.lock().await.is_empty() {
while metadata_backend.state.lock().await.objects.is_empty() {
tokio::task::yield_now().await;
}
})
@@ -2172,7 +2173,7 @@ async fn create_table_holds_bucket_fence_from_metadata_write_through_registratio
#[tokio::test]
async fn create_table_response_recreates_dropped_identifier_without_overwriting_retained_metadata() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -2308,9 +2309,9 @@ async fn create_table_response_recreates_dropped_identifier_without_overwriting_
#[tokio::test]
async fn concurrent_create_table_responses_keep_one_catalog_winner_with_distinct_metadata() {
let catalog_backend = TestTableCatalogObjectBackend::default();
let catalog_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(catalog_backend);
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -2404,7 +2405,7 @@ async fn concurrent_create_table_responses_keep_one_catalog_winner_with_distinct
#[tokio::test]
async fn standard_commit_applies_updates_and_writes_next_metadata() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_uuid = created.metadata["table-uuid"]
@@ -2535,7 +2536,7 @@ fn table_metadata_file_name_scoping_is_bounded_and_identity_sensitive() {
#[tokio::test]
async fn renamed_and_recreated_tables_with_the_same_commit_id_use_disjoint_metadata_files() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::StrongTableCatalogStore::new(metadata_backend.clone());
let source_namespace = crate::table_catalog::Namespace::parse("analytics").expect("source namespace should parse");
let destination_namespace = crate::table_catalog::Namespace::parse("curated").expect("destination namespace should parse");
@@ -2952,7 +2953,7 @@ async fn standard_commit_recovers_matching_table_scoped_metadata_orphan() {
#[tokio::test]
async fn concurrent_identical_commits_reuse_table_scoped_metadata_winner() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::StrongTableCatalogStore::new(metadata_backend.clone());
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -3192,7 +3193,7 @@ async fn standard_commit_rejects_fallback_readback_mismatch() {
#[tokio::test]
async fn standard_commit_uses_client_uuid_commit_id_in_metadata_file_name() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -3241,7 +3242,7 @@ async fn standard_commit_uses_client_uuid_commit_id_in_metadata_file_name() {
#[tokio::test]
async fn standard_commit_accepts_non_uuid_client_commit_id_without_using_it_in_metadata_file_name() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -3284,7 +3285,7 @@ async fn standard_commit_accepts_non_uuid_client_commit_id_without_using_it_in_m
#[tokio::test]
async fn commit_publication_uses_idempotency_key_as_retry_identity() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
let idempotency_key = Uuid::now_v7().to_string();
@@ -3319,7 +3320,7 @@ async fn commit_publication_replays_historical_standard_commit_across_backings()
crate::table_catalog::TableCatalogBackingMode::ObjectBacked,
crate::table_catalog::TableCatalogBackingMode::DurableStrong,
] {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ConfiguredTableCatalogStore::new_for_test(metadata_backend.clone(), mode);
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -3435,7 +3436,7 @@ async fn commit_publication_replays_historical_standard_commit_across_backings()
#[tokio::test]
async fn staged_standard_commit_retry_revalidates_referenced_objects() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -3545,7 +3546,7 @@ async fn staged_standard_commit_retry_revalidates_referenced_objects() {
#[tokio::test]
async fn commit_publication_denies_generated_metadata_write_before_pointer_advance() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
let before = store
@@ -3605,7 +3606,7 @@ async fn commit_publication_denies_generated_metadata_write_before_pointer_advan
#[tokio::test]
async fn commit_publication_authorizes_referenced_objects() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let entry = store
@@ -3672,7 +3673,7 @@ async fn commit_publication_authorizes_referenced_objects() {
#[tokio::test]
async fn commit_publication_denies_referenced_data_read_before_pointer_advance() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let before = store
@@ -3738,7 +3739,7 @@ async fn commit_publication_holds_referenced_object_locks_until_pointer_publish(
commit_table_pause: Some(pause.clone()),
..Default::default()
});
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(store.as_ref(), &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -3846,7 +3847,7 @@ async fn commit_publication_holds_referenced_object_locks_until_pointer_publish(
#[tokio::test]
async fn rolling_upgrade_commit_retains_legacy_data_file_guard_until_publication_completes() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let data_file = "tables/table-id/data/part-00001.parquet";
metadata_backend.put_bytes("warehouse", data_file, b"data".to_vec()).await;
let commit_backend = TableCommitObjectBackend::rolling_upgrade(metadata_backend.clone());
@@ -3908,7 +3909,7 @@ async fn rolling_upgrade_initial_publication_fences_old_and_new_data_plane_write
.expect("namespace should seed");
let data_file = "tables/table-id/data/part-00001.parquet";
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
metadata_backend.put_bytes("warehouse", data_file, b"data".to_vec()).await;
let publication_backend = TableCommitObjectBackend::rolling_upgrade(metadata_backend.clone());
assert!(
@@ -4004,7 +4005,7 @@ async fn warehouse_relocation_holds_bucket_fence_before_catalog_publication() {
commit_table_pause: Some(pause.clone()),
..Default::default()
});
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(store.as_ref(), &metadata_backend, &namespace).await;
let current = store
@@ -4079,7 +4080,7 @@ async fn warehouse_relocation_holds_bucket_fence_before_catalog_publication() {
#[tokio::test]
async fn commit_publication_lock_order_remains_compatible_with_old_maintenance_nodes() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let table = crate::table_catalog::IdentifierSegment::parse("events").expect("table should parse");
let current_metadata = crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00001.metadata.json");
@@ -4163,7 +4164,7 @@ async fn commit_publication_lock_order_remains_compatible_with_old_maintenance_n
#[tokio::test]
async fn commit_publication_acquires_discovered_object_locks_in_key_order() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let first = "metadata/a.json";
let last = "metadata/z.json";
metadata_backend.put_bytes("warehouse", first, b"a".to_vec()).await;
@@ -4238,7 +4239,7 @@ async fn commit_publication_acquires_discovered_object_locks_in_key_order() {
#[tokio::test]
async fn commit_publication_revalidates_objects_after_ordered_lock_acquisition() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let object = "metadata/current.json";
metadata_backend.put_bytes("warehouse", object, b"before".to_vec()).await;
let commit_backend = TableCommitObjectBackend::trusted(metadata_backend.clone());
@@ -4260,16 +4261,16 @@ async fn commit_publication_revalidates_objects_after_ordered_lock_acquisition()
#[tokio::test]
async fn commit_publication_binds_fingerprint_to_returned_bytes() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let object = "metadata/current.json";
let original = b"original".to_vec();
let replacement = b"replacement".to_vec();
let original_etag = hex_sha256(&original, str::to_string);
metadata_backend.objects.lock().await.insert(
metadata_backend.state.lock().await.objects.insert(
("warehouse".to_string(), object.to_string()),
crate::table_catalog::TableCatalogObject {
TestCatalogObjectRecord {
data: replacement,
etag: Some(original_etag.clone()),
etag: original_etag.clone(),
mod_time: None,
},
);
@@ -4278,11 +4279,11 @@ async fn commit_publication_binds_fingerprint_to_returned_bytes() {
.await
.expect("replacement bytes should be discovered")
.expect("replacement object should exist");
metadata_backend.objects.lock().await.insert(
metadata_backend.state.lock().await.objects.insert(
("warehouse".to_string(), object.to_string()),
crate::table_catalog::TableCatalogObject {
TestCatalogObjectRecord {
data: original,
etag: Some(original_etag),
etag: original_etag,
mod_time: None,
},
);
@@ -4302,7 +4303,7 @@ async fn commit_publication_binds_fingerprint_to_returned_bytes() {
#[tokio::test]
async fn standard_commit_publishes_more_than_ten_thousand_live_files() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -4330,15 +4331,15 @@ async fn standard_commit_publishes_more_than_ten_thousand_live_files() {
)
.await;
{
let mut objects = metadata_backend.objects.lock().await;
let mut state = metadata_backend.state.lock().await;
let data = vec![1];
let etag = hex_sha256(&data, str::to_string);
for file in &data_files {
objects.insert(
state.objects.insert(
("warehouse".to_string(), test_snapshot_object_key("warehouse", file)),
crate::table_catalog::TableCatalogObject {
TestCatalogObjectRecord {
data: data.clone(),
etag: Some(etag.clone()),
etag: etag.clone(),
mod_time: None,
},
);
@@ -4390,7 +4391,7 @@ async fn standard_commit_publishes_more_than_ten_thousand_live_files() {
#[tokio::test]
async fn commit_publication_rejects_recreated_object_observed_by_exists() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let object = "data/part-00001.parquet";
metadata_backend.put_bytes("warehouse", object, b"before".to_vec()).await;
let commit_backend = TableCommitObjectBackend::trusted(metadata_backend.clone());
@@ -4419,7 +4420,7 @@ async fn commit_publication_rejects_recreated_object_observed_by_exists() {
#[tokio::test]
async fn standard_commit_ignores_generation_only_orphan_metadata_file() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
metadata_backend
@@ -4468,15 +4469,15 @@ async fn standard_commit_ignores_generation_only_orphan_metadata_file() {
#[tokio::test]
async fn concurrent_standard_commits_write_distinct_metadata_files_before_pointer_conflict() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
let barrier = Arc::new(tokio::sync::Barrier::new(2));
let metadata_backend = TestTableCatalogObjectBackend {
objects: Arc::clone(&metadata_backend.objects),
state: Arc::clone(&metadata_backend.state),
put_object_barrier: Some(barrier),
..Default::default()
..TestTableCatalogObjectBackend::content_addressed()
};
let first_commit_id = "33333333-3333-4333-8333-333333333333";
let second_commit_id = "44444444-4444-4444-8444-444444444444";
@@ -4537,7 +4538,7 @@ async fn concurrent_standard_commits_write_distinct_metadata_files_before_pointe
#[tokio::test]
async fn standard_commit_accepts_legacy_catalog_uuid_when_current_metadata_matches() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -4607,7 +4608,7 @@ async fn standard_commit_accepts_legacy_catalog_uuid_when_current_metadata_match
#[tokio::test]
async fn metadata_location_api_accepts_legacy_catalog_uuid_when_target_matches_current_metadata() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -4673,7 +4674,7 @@ async fn metadata_location_api_accepts_legacy_catalog_uuid_when_target_matches_c
#[tokio::test]
async fn table_metadata_maintenance_helper_runs_dry_run_and_delete() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -4847,7 +4848,7 @@ async fn table_metadata_maintenance_helper_runs_dry_run_and_delete() {
#[tokio::test]
async fn table_metadata_maintenance_helper_commits_snapshot_expiration() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5003,7 +5004,7 @@ async fn table_metadata_maintenance_helper_commits_snapshot_expiration() {
#[tokio::test]
async fn table_metadata_maintenance_helper_commits_compaction_through_publication_observer() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5095,7 +5096,7 @@ async fn table_metadata_maintenance_helper_commits_compaction_through_publicatio
#[tokio::test]
async fn table_metadata_maintenance_helper_rejects_snapshot_expiration_manual_review_commit() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5173,7 +5174,7 @@ async fn table_metadata_maintenance_helper_rejects_snapshot_expiration_manual_re
#[tokio::test]
async fn table_metadata_maintenance_helper_rejects_stale_snapshot_expiration_plan() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5257,7 +5258,7 @@ async fn table_metadata_maintenance_helper_rejects_stale_snapshot_expiration_pla
#[tokio::test]
async fn table_metadata_maintenance_helper_rejects_delete_with_snapshot_expiration_commit() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5286,7 +5287,7 @@ async fn table_metadata_maintenance_helper_rejects_delete_with_snapshot_expirati
#[tokio::test]
async fn table_refs_response_reports_current_and_user_defined_refs() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5327,7 +5328,7 @@ async fn table_refs_response_reports_current_and_user_defined_refs() {
#[tokio::test]
async fn external_catalog_bridge_response_lists_supported_operator_bridges() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend);
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5362,7 +5363,7 @@ async fn external_catalog_bridge_response_lists_supported_operator_bridges() {
#[tokio::test]
async fn external_catalog_bridge_persists_identity_and_boundary() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend);
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5428,7 +5429,7 @@ async fn external_catalog_bridge_persists_identity_and_boundary() {
#[tokio::test]
async fn external_catalog_bridge_sync_registers_missing_table_from_snapshot() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5506,7 +5507,7 @@ async fn external_catalog_bridge_sync_registers_missing_table_from_snapshot() {
#[tokio::test]
async fn external_catalog_bridge_sync_commits_existing_table_pointer() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5563,7 +5564,7 @@ async fn external_catalog_bridge_sync_commits_existing_table_pointer() {
#[tokio::test]
async fn external_catalog_bridge_sync_denies_metadata_reads_before_pointer_publish() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5632,7 +5633,7 @@ async fn external_catalog_bridge_sync_denies_metadata_reads_before_pointer_publi
#[tokio::test]
async fn external_catalog_bridge_sync_conflicts_leave_pointer_unchanged() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -5823,7 +5824,7 @@ fn snapshot_conflict_rejects_unknown_snapshot_operations() {
#[tokio::test]
async fn row_level_conflict_allows_overwrite_when_deleted_file_is_current() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -5959,7 +5960,7 @@ async fn row_level_conflict_allows_overwrite_when_deleted_file_is_current() {
#[tokio::test]
async fn row_level_conflict_allows_v1_manifest_snapshot() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6069,7 +6070,7 @@ async fn row_level_conflict_allows_v1_manifest_snapshot() {
#[tokio::test]
async fn row_level_conflict_inherits_manifest_list_sequence_numbers() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6122,7 +6123,7 @@ async fn row_level_conflict_inherits_manifest_list_sequence_numbers() {
#[tokio::test]
async fn row_level_conflict_allows_inherited_manifests_on_append() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6228,7 +6229,7 @@ async fn row_level_conflict_allows_inherited_manifests_on_append() {
#[tokio::test]
async fn row_level_conflict_rejects_changed_inherited_manifest_identity() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6330,7 +6331,7 @@ async fn row_level_conflict_rejects_changed_inherited_manifest_identity() {
#[tokio::test]
async fn row_level_conflict_rejects_stale_new_manifest_sequence() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6389,7 +6390,7 @@ async fn row_level_conflict_rejects_stale_new_manifest_sequence() {
#[tokio::test]
async fn row_level_conflict_rejects_stale_added_entry_sequence() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6448,7 +6449,7 @@ async fn row_level_conflict_rejects_stale_added_entry_sequence() {
#[tokio::test]
async fn row_level_conflict_rejects_historical_change_in_new_manifest() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6507,7 +6508,7 @@ async fn row_level_conflict_rejects_historical_change_in_new_manifest() {
#[tokio::test]
async fn row_level_conflict_allows_add_only_overwrite_snapshot() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6618,7 +6619,7 @@ async fn row_level_conflict_allows_add_only_overwrite_snapshot() {
#[tokio::test]
async fn row_level_conflict_rejects_delete_of_non_current_file() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6737,7 +6738,7 @@ async fn row_level_conflict_rejects_delete_of_non_current_file() {
#[tokio::test]
async fn row_level_conflict_rejects_append_with_delete_files() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6794,7 +6795,7 @@ async fn row_level_conflict_rejects_append_with_delete_files() {
#[tokio::test]
async fn row_level_conflict_rejects_missing_manifest_before_pointer_update() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -6900,7 +6901,7 @@ async fn row_level_conflict_rejects_missing_manifest_before_pointer_update() {
#[tokio::test]
async fn row_level_conflict_rejects_manifest_outside_table_warehouse() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -7133,7 +7134,7 @@ fn create_view_request_accepts_deep_warehouse_location() {
#[tokio::test]
async fn view_catalog_responses_persist_replace_and_drop_view_metadata() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -7297,7 +7298,7 @@ async fn view_catalog_responses_persist_replace_and_drop_view_metadata() {
#[tokio::test]
async fn table_ref_write_responses_use_commit_guard_and_protect_deletes() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
@@ -7949,90 +7950,6 @@ impl TestCatalogPublishPause {
}
}
type TestTableCatalogObjectLocks = Arc<tokio::sync::Mutex<BTreeMap<(String, String), Arc<tokio::sync::Mutex<()>>>>>;
#[derive(Clone, Default)]
struct TestTableCatalogObjectBackend {
objects: Arc<tokio::sync::Mutex<BTreeMap<(String, String), crate::table_catalog::TableCatalogObject>>>,
put_object_barrier: Option<Arc<tokio::sync::Barrier>>,
fail_put_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
corrupt_put_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
missing_read_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
fail_read_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
locks: TestTableCatalogObjectLocks,
lock_attempts: Arc<tokio::sync::Mutex<Vec<(String, String)>>>,
}
impl TestTableCatalogObjectBackend {
async fn put_bytes(&self, bucket: &str, object: &str, data: Vec<u8>) {
let etag = hex_sha256(&data, str::to_string);
self.objects.lock().await.insert(
(bucket.to_string(), object.to_string()),
crate::table_catalog::TableCatalogObject {
data,
etag: Some(etag),
mod_time: None,
},
);
}
async fn put_json(&self, bucket: &str, object: &str, value: serde_json::Value) {
self.put_json_with_mod_time(bucket, object, value, None).await;
}
async fn put_gzip_json(&self, bucket: &str, object: &str, value: serde_json::Value) {
use std::io::Write;
let data = serde_json::to_vec(&value).expect("metadata JSON should serialize");
let mut encoder = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
encoder.write_all(&data).expect("metadata JSON should compress");
self.put_bytes(bucket, object, encoder.finish().expect("metadata gzip stream should finish"))
.await;
}
async fn put_json_with_mod_time(
&self,
bucket: &str,
object: &str,
value: serde_json::Value,
mod_time: Option<OffsetDateTime>,
) {
let data = serde_json::to_vec(&value).expect("metadata JSON should serialize");
let etag = hex_sha256(&data, str::to_string);
self.objects.lock().await.insert(
(bucket.to_string(), object.to_string()),
crate::table_catalog::TableCatalogObject {
data,
etag: Some(etag),
mod_time,
},
);
}
async fn write_lock_is_held(&self, bucket: &str, object: &str) -> bool {
let lock = self
.locks
.lock()
.await
.get(&(bucket.to_string(), object.to_string()))
.cloned();
lock.is_some_and(|lock| lock.try_lock_owned().is_err())
}
async fn wait_for_lock_attempts(&self, count: usize) {
tokio::time::timeout(StdDuration::from_secs(2), async {
loop {
if self.lock_attempts.lock().await.len() >= count {
return;
}
tokio::task::yield_now().await;
}
})
.await
.expect("lock acquisition attempts should be observable");
}
}
fn trusted_table_commit_backend(
backend: &TestTableCatalogObjectBackend,
) -> TableCommitObjectBackend<TestTableCatalogObjectBackend> {
@@ -8210,7 +8127,7 @@ async fn standard_commit_foreign_primary_fixture() -> (
String,
) {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
let current = store
@@ -8252,7 +8169,7 @@ async fn standard_commit_primary_fixture(
String,
) {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_standard_events_table(&store, &metadata_backend, &namespace).await;
let current = store
@@ -8401,130 +8318,6 @@ async fn seed_object_table_for_metadata_maintenance(
.await;
}
#[async_trait::async_trait]
impl crate::table_catalog::TableCatalogObjectBackend for TestTableCatalogObjectBackend {
async fn read_object(
&self,
bucket: &str,
object: &str,
) -> crate::table_catalog::TableCatalogStoreResult<Option<crate::table_catalog::TableCatalogObject>> {
let mut missing_read_object_path = self.missing_read_object_path.lock().await;
if missing_read_object_path.as_deref() == Some(object) {
missing_read_object_path.take();
return Ok(None);
}
drop(missing_read_object_path);
let mut fail_read_object_path = self.fail_read_object_path.lock().await;
if fail_read_object_path.as_deref() == Some(object) {
fail_read_object_path.take();
return Err(crate::table_catalog::TableCatalogStoreError::Internal(
"private generated metadata read failure".to_string(),
));
}
drop(fail_read_object_path);
Ok(self
.objects
.lock()
.await
.get(&(bucket.to_string(), object.to_string()))
.cloned())
}
async fn object_exists(&self, bucket: &str, object: &str) -> crate::table_catalog::TableCatalogStoreResult<bool> {
Ok(self
.objects
.lock()
.await
.contains_key(&(bucket.to_string(), object.to_string())))
}
async fn put_object(
&self,
bucket: &str,
object: &str,
data: Vec<u8>,
precondition: crate::table_catalog::TableCatalogPutPrecondition,
) -> crate::table_catalog::TableCatalogStoreResult<()> {
let mut fail_put_object_path = self.fail_put_object_path.lock().await;
if fail_put_object_path.as_deref() == Some(object) {
fail_put_object_path.take();
return Err(crate::table_catalog::TableCatalogStoreError::Internal(
"injected metadata write failure".to_string(),
));
}
drop(fail_put_object_path);
let mut corrupt_put_object_path = self.corrupt_put_object_path.lock().await;
let data = if corrupt_put_object_path.as_deref() == Some(object) {
corrupt_put_object_path.take();
b"{}".to_vec()
} else {
data
};
drop(corrupt_put_object_path);
let key = (bucket.to_string(), object.to_string());
let mut objects = self.objects.lock().await;
let result = if matches!(precondition, crate::table_catalog::TableCatalogPutPrecondition::IfAbsent)
&& objects.contains_key(&key)
{
Err(crate::table_catalog::TableCatalogStoreError::Conflict(format!(
"object already exists: {object}"
)))
} else {
let etag = hex_sha256(&data, str::to_string);
objects.insert(
key,
crate::table_catalog::TableCatalogObject {
data,
etag: Some(etag),
mod_time: None,
},
);
Ok(())
};
drop(objects);
if let Some(barrier) = &self.put_object_barrier {
barrier.wait().await;
}
result
}
async fn delete_object(&self, bucket: &str, object: &str) -> crate::table_catalog::TableCatalogStoreResult<()> {
self.objects.lock().await.remove(&(bucket.to_string(), object.to_string()));
Ok(())
}
async fn list_objects(&self, bucket: &str, prefix: &str) -> crate::table_catalog::TableCatalogStoreResult<Vec<String>> {
Ok(self
.objects
.lock()
.await
.keys()
.filter(|(object_bucket, object)| object_bucket == bucket && object.starts_with(prefix))
.map(|(_, object)| object.clone())
.collect())
}
async fn acquire_write_lock(
&self,
bucket: &str,
object: &str,
) -> crate::table_catalog::TableCatalogStoreResult<Box<dyn Send>> {
self.lock_attempts.lock().await.push((bucket.to_string(), object.to_string()));
let lock = {
let mut locks = self.locks.lock().await;
locks
.entry((bucket.to_string(), object.to_string()))
.or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
.clone()
};
Ok(Box::new(lock.lock_owned().await))
}
}
#[async_trait::async_trait]
impl crate::table_catalog::TableCatalogStore for TestTableCatalogStore {
async fn get_table_bucket(
@@ -9040,7 +8833,7 @@ async fn namespace_helpers_call_catalog_store() {
#[tokio::test]
async fn table_helpers_call_catalog_store() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -9172,7 +8965,7 @@ async fn table_helpers_call_catalog_store() {
#[tokio::test]
async fn register_table_response_adopts_metadata_table_uuid() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -9223,7 +9016,7 @@ async fn register_table_response_adopts_metadata_table_uuid() {
#[tokio::test]
async fn register_table_denies_metadata_read_before_catalog_publication() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let (namespace, metadata_location) = seed_events_registration_target(&store, &metadata_backend).await;
let authorized = Arc::new(tokio::sync::Mutex::new(Vec::new()));
let guarded_backend =
@@ -9266,7 +9059,7 @@ async fn register_table_denies_metadata_read_before_catalog_publication() {
#[tokio::test]
async fn catalog_import_denies_metadata_read_before_catalog_publication() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let (namespace, metadata_location) = seed_events_registration_target(&store, &metadata_backend).await;
let authorized = Arc::new(tokio::sync::Mutex::new(Vec::new()));
let guarded_backend =
@@ -9308,7 +9101,7 @@ async fn catalog_import_denies_metadata_read_before_catalog_publication() {
#[tokio::test]
async fn register_table_rejects_metadata_replaced_before_catalog_publication() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = Arc::new(crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()));
let (namespace, metadata_location) = seed_events_registration_target(store.as_ref(), &metadata_backend).await;
let table = crate::table_catalog::IdentifierSegment::parse("events").expect("table should parse");
@@ -9376,7 +9169,7 @@ async fn register_table_rejects_metadata_replaced_before_catalog_publication() {
#[tokio::test]
async fn catalog_import_rejects_metadata_replaced_before_catalog_publication() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = Arc::new(crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()));
let (namespace, metadata_location) = seed_events_registration_target(store.as_ref(), &metadata_backend).await;
let table = crate::table_catalog::IdentifierSegment::parse("events").expect("table should parse");
@@ -9445,7 +9238,7 @@ async fn catalog_import_rejects_metadata_replaced_before_catalog_publication() {
#[tokio::test]
async fn register_table_response_rejects_metadata_without_format_version() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -9501,7 +9294,7 @@ async fn register_table_response_rejects_metadata_without_format_version() {
#[tokio::test]
async fn metadata_location_api_loads_and_updates_current_pointer() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -9577,7 +9370,7 @@ async fn metadata_location_api_loads_and_updates_current_pointer() {
#[tokio::test]
async fn metadata_location_api_accepts_gzip_table_metadata() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let current = store
@@ -9616,7 +9409,7 @@ async fn metadata_location_api_accepts_gzip_table_metadata() {
#[tokio::test]
async fn metadata_location_api_validates_snapshot_graph_before_commit() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let current = store
@@ -9682,7 +9475,7 @@ async fn metadata_location_api_validates_snapshot_graph_before_commit() {
#[tokio::test]
async fn metadata_location_api_validates_relocated_snapshot_graph_under_target_warehouse() {
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
@@ -9741,7 +9534,7 @@ async fn metadata_location_api_validates_relocated_snapshot_graph_under_target_w
#[tokio::test]
async fn metadata_location_api_rejects_invalid_target_metadata_before_commit() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -9819,7 +9612,7 @@ async fn metadata_location_api_rejects_invalid_target_metadata_before_commit() {
#[tokio::test]
async fn metadata_location_api_rejects_mismatched_table_uuid_before_commit() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -9898,7 +9691,7 @@ async fn metadata_location_api_rejects_mismatched_table_uuid_before_commit() {
#[tokio::test]
async fn catalog_import_and_rollback_use_register_and_commit_paths() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -9993,7 +9786,7 @@ async fn catalog_import_and_rollback_use_register_and_commit_paths() {
#[tokio::test]
async fn rollback_denies_metadata_reads_before_pointer_publish() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -10055,7 +9848,7 @@ async fn rollback_denies_metadata_reads_before_pointer_publish() {
#[tokio::test]
async fn rollback_rejects_invalid_target_metadata_before_commit() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -10144,7 +9937,7 @@ async fn rollback_rejects_invalid_target_metadata_before_commit() {
#[tokio::test]
async fn rollback_rejects_mismatched_table_uuid_before_commit() {
let backend = TestTableCatalogObjectBackend::default();
let backend = TestTableCatalogObjectBackend::content_addressed();
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
let bucket = "warehouse";
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
@@ -10229,7 +10022,7 @@ async fn rollback_rejects_mismatched_table_uuid_before_commit() {
#[tokio::test]
async fn legacy_commit_rejects_mismatched_table_uuid_before_commit() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let metadata_backend = TestTableCatalogObjectBackend::content_addressed();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
+24 -42
View File
@@ -86,7 +86,7 @@ use super::storage_api::object_usecase::options::{
namespace_reserved_user_metadata, normalize_content_encoding_for_storage, preserve_unclassified_user_metadata,
put_opts_with_replication_authorization, validate_archive_content_encoding,
};
use super::storage_api::object_usecase::request_context::{self, spawn_traced, spawn_traced_join};
use super::storage_api::object_usecase::request_context::{self, spawn_traced};
use super::storage_api::object_usecase::s3_api::multipart::parse_list_parts_params;
use super::storage_api::object_usecase::set_disk::{
get_lock_acquire_timeout, get_object_disk_read_timeout, is_valid_storage_class,
@@ -7506,50 +7506,31 @@ impl DefaultObjectUsecase {
let cache_adapter = self.object_data_cache();
let _ = invalidate_object_data_cache_before_mutation(&cache_adapter, &bucket, &key).await;
let copy_commit = spawn_traced_join({
let store = Arc::clone(&store);
let src_bucket = src_bucket.clone();
let src_key = src_key.clone();
let bucket = bucket.clone();
let key = key.clone();
let src_opts = src_opts.clone();
let dst_opts = dst_opts.clone();
async move {
let _source_bucket_lifecycle_guard = source_bucket_lifecycle_guard;
let _destination_bucket_lifecycle_guard_storage = destination_bucket_lifecycle_guard_storage;
let _self_copy_lock_guard = _self_copy_lock_guard;
let oi = store
.copy_object(&src_bucket, &src_key, &bucket, &key, &mut src_info, &src_opts, &dst_opts)
.await
.map_err(ApiError::from)?;
drop(_self_copy_lock_guard);
let oi = store
.copy_object(&src_bucket, &src_key, &bucket, &key, &mut src_info, &src_opts, &dst_opts)
.await
.map_err(ApiError::from)?;
// Reuse the single pre-commit replication decision (see `dsc` above) so
// the persisted pending marker and the schedule always agree, mirroring
// the PUT path.
if dsc.replicate_any() {
schedule_object_replication(oi.clone(), store.clone(), dsc).await;
}
// Reuse the single pre-commit replication decision (see `dsc` above) so
// the persisted pending marker and the schedule always agree, mirroring
// the PUT path.
if dsc.replicate_any() {
schedule_object_replication(oi.clone(), Arc::clone(&store), dsc).await;
}
maybe_enqueue_transition_immediate(&oi, LcEventSrc::S3CopyObject).await;
let _ = invalidate_object_data_cache_after_copy_success(&cache_adapter, &bucket, &key).await;
maybe_enqueue_transition_immediate(&oi, LcEventSrc::S3CopyObject).await;
let _ = invalidate_object_data_cache_after_copy_success(&cache_adapter, &bucket, &key).await;
let dest_versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
if has_bucket_metadata {
if dest_versioned {
record_bucket_object_version_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await;
} else {
record_bucket_object_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await;
}
}
rustfs_scanner::record_dirty_usage_bucket(&bucket);
Ok::<_, S3Error>((oi, dest_versioned))
let dest_versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
// Update quota tracking after successful copy
if has_bucket_metadata {
if dest_versioned {
record_bucket_object_version_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await;
} else {
record_bucket_object_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await;
}
});
let (oi, dest_versioned) = copy_commit.await.map_err(|err| {
S3Error::with_message(S3ErrorCode::InternalError, format!("copy object commit owner task failed: {err}"))
})??;
}
let raw_dest_version = oi.version_id.map(|v| v.to_string());
let dest_version = if dest_versioned { raw_dest_version } else { None };
@@ -7597,7 +7578,7 @@ impl DefaultObjectUsecase {
}
}
let copy_object_result = CopyObjectResult {
e_tag: oi.etag.as_ref().map(|etag| to_s3s_etag(etag)),
e_tag: oi.etag.map(|etag| to_s3s_etag(&etag)),
last_modified: oi.mod_time.map(Timestamp::from),
checksum_crc32: response_checksums.crc32,
checksum_crc32c: response_checksums.crc32c,
@@ -7628,6 +7609,7 @@ impl DefaultObjectUsecase {
let result = Ok(S3Response::new(output));
let _ = helper.complete(&result);
rustfs_scanner::record_dirty_usage_bucket(&bucket);
result
}
+1 -1
View File
@@ -998,7 +998,7 @@ pub(crate) mod options {
}
pub(crate) mod request_context {
pub(crate) use crate::storage::storage_api::request_context_consumer::{RequestContext, spawn_traced, spawn_traced_join};
pub(crate) use crate::storage::storage_api::request_context_consumer::{RequestContext, spawn_traced};
}
pub(crate) mod sse {
-9
View File
@@ -257,15 +257,6 @@ where
tokio::spawn(tracing::Instrument::instrument(fut, tracing::Span::current()));
}
/// Spawn a request-internal task and return its join handle to the caller.
pub fn spawn_traced_join<F>(fut: F) -> tokio::task::JoinHandle<F::Output>
where
F: std::future::Future + Send + 'static,
F::Output: Send + 'static,
{
tokio::spawn(tracing::Instrument::instrument(fut, tracing::Span::current()))
}
#[cfg(test)]
#[allow(unused_imports)]
mod tests {
+1 -3
View File
@@ -203,9 +203,7 @@ pub(crate) mod options_consumer {
}
pub(crate) mod request_context_consumer {
pub(crate) use super::super::request_context::{
RequestContext, extract_request_id_from_headers, spawn_traced, spawn_traced_join,
};
pub(crate) use super::super::request_context::{RequestContext, extract_request_id_from_headers, spawn_traced};
}
pub(crate) mod rpc_consumer {
+713
View File
@@ -22,6 +22,15 @@
//! fixed values (sequence 7 / snapshot 20), which keeps every produced byte
//! identical to the pre-extraction fixtures.
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use time::OffsetDateTime;
use super::{
StrongTableCatalogRuntime, TableCatalogObject, TableCatalogObjectBackend, TableCatalogObjectMetadata,
TableCatalogPutPrecondition, TableCatalogStoreError, TableCatalogStoreResult, TableCommitPublication,
};
pub(crate) fn table_metadata_json(table_uuid: &str, location: &str) -> serde_json::Value {
serde_json::json!({
"format-version": 2,
@@ -226,3 +235,707 @@ pub(crate) fn manifest_avro_bytes_with_nullable_sequences(files: &[(&str, i32, i
}
writer.into_inner().expect("manifest avro bytes should flush")
}
// --- Stateful object backend shared by the store and admin handler tests
// (backlog#1837 PR2). Superset instrumentation lands here incrementally;
// this is the store-side fake moved verbatim.
#[derive(Clone, Default)]
pub(crate) struct TestCatalogObjectBackend {
pub(crate) state: Arc<tokio::sync::Mutex<TestCatalogObjectState>>,
pub(crate) locks: TestCatalogObjectLocks,
pub(crate) strong_runtime: Option<StrongTableCatalogRuntime>,
// One-shot, path-keyed injection knobs from the admin handler tests'
// former TestTableCatalogObjectBackend (backlog#1837 PR2): each fires
// once for the named object and clears itself, mirroring the original
// semantics exactly. They compose with (and run before) the store tests'
// attempt-indexed injection maps above.
pub(crate) put_object_barrier: Option<Arc<tokio::sync::Barrier>>,
pub(crate) fail_put_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
pub(crate) corrupt_put_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
pub(crate) missing_read_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
pub(crate) fail_read_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
pub(crate) lock_attempts: Arc<tokio::sync::Mutex<Vec<(String, String)>>>,
/// Content-addressed (sha256) etags instead of the store fake's counter.
/// The admin handler tests observe an object's etag and expect rewriting
/// identical bytes to reproduce it, so their fixtures set this.
pub(crate) content_addressed_etags: bool,
}
pub(crate) type TestCatalogObjectLockKey = (String, String);
pub(crate) type TestCatalogObjectLock = Arc<tokio::sync::RwLock<()>>;
pub(crate) type TestCatalogObjectLocks = Arc<tokio::sync::Mutex<BTreeMap<TestCatalogObjectLockKey, TestCatalogObjectLock>>>;
#[derive(Clone, Default)]
pub(crate) struct TestCatalogObjectPause {
started: Arc<tokio::sync::Notify>,
release: Arc<tokio::sync::Notify>,
}
impl TestCatalogObjectPause {
pub(crate) async fn wait_started(&self) {
self.started.notified().await;
}
pub(crate) fn release(&self) {
self.release.notify_one();
}
}
#[derive(Clone)]
pub(crate) struct BlockingObjectPublication {
backend: TestCatalogObjectBackend,
object: String,
started: Arc<tokio::sync::Notify>,
guard: Arc<parking_lot::Mutex<Option<Box<dyn Send>>>>,
}
impl BlockingObjectPublication {
pub(crate) fn new(backend: TestCatalogObjectBackend, object: impl Into<String>) -> Self {
Self {
backend,
object: object.into(),
started: Arc::new(tokio::sync::Notify::new()),
guard: Arc::new(parking_lot::Mutex::new(None)),
}
}
pub(crate) async fn wait_started(&self) {
self.started.notified().await;
}
}
#[derive(Default)]
pub(crate) struct UnserializedTestPublication;
#[async_trait::async_trait]
impl TableCommitPublication for UnserializedTestPublication {
async fn begin_table_bucket(&self, _table_bucket: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
async fn prepare(&self, _table_bucket: &str, _namespace: &str, _table: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
fn holds_table_bucket(&self, _table_bucket: &str) -> bool {
true
}
fn holds_table(&self, _table_bucket: &str, _namespace: &str, _table: &str) -> bool {
true
}
fn complete(&self) {}
}
#[async_trait::async_trait]
impl TableCommitPublication for BlockingObjectPublication {
async fn begin_table_bucket(&self, _table_bucket: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
async fn prepare(&self, table_bucket: &str, _namespace: &str, _table: &str) -> TableCatalogStoreResult<()> {
self.started.notify_one();
let guard = self.backend.acquire_read_lock(table_bucket, &self.object).await?;
*self.guard.lock() = Some(guard);
Ok(())
}
fn holds_table_bucket(&self, _table_bucket: &str) -> bool {
true
}
fn holds_table(&self, _table_bucket: &str, _namespace: &str, _table: &str) -> bool {
self.guard.lock().is_some()
}
fn complete(&self) {
drop(self.guard.lock().take());
}
}
#[derive(Default)]
pub(crate) struct TestCatalogObjectState {
pub(crate) objects: BTreeMap<(String, String), TestCatalogObjectRecord>,
pub(crate) etagless_objects: BTreeSet<(String, String)>,
pub(crate) fail_read_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pub(crate) pause_before_read_attempts: BTreeMap<(String, String), BTreeMap<usize, TestCatalogObjectPause>>,
pub(crate) pause_read_attempts: BTreeMap<(String, String), BTreeMap<usize, TestCatalogObjectPause>>,
pub(crate) read_attempts: BTreeMap<(String, String), usize>,
pub(crate) read_limits: Vec<((String, String), usize)>,
pub(crate) fail_put_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pub(crate) fail_after_put_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pub(crate) pause_put_attempts: BTreeMap<(String, String), BTreeMap<usize, TestCatalogObjectPause>>,
pub(crate) fail_delete_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pub(crate) fail_after_delete_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pub(crate) put_attempts: BTreeMap<(String, String), usize>,
pub(crate) delete_attempts: BTreeMap<(String, String), usize>,
pub(crate) write_lock_acquisitions: BTreeMap<(String, String), usize>,
pub(crate) read_lock_acquisitions: BTreeMap<(String, String), usize>,
pub(crate) read_calls: usize,
pub(crate) metadata_calls: usize,
pub(crate) list_calls: usize,
pub(crate) next_etag: u64,
}
#[derive(Clone)]
pub(crate) struct TestCatalogObjectRecord {
pub(crate) data: Vec<u8>,
pub(crate) etag: String,
pub(crate) mod_time: Option<OffsetDateTime>,
}
impl TestCatalogObjectBackend {
pub(crate) async fn seed_object(&self, bucket: &str, object: &str, data: Vec<u8>) {
self.seed_object_with_mod_time(bucket, object, data, Some(OffsetDateTime::UNIX_EPOCH))
.await;
}
pub(crate) async fn seed_object_with_mod_time(
&self,
bucket: &str,
object: &str,
data: Vec<u8>,
mod_time: Option<OffsetDateTime>,
) {
let mut state = self.state.lock().await;
let etag = state.next_etag();
state
.objects
.insert((bucket.to_string(), object.to_string()), TestCatalogObjectRecord { data, etag, mod_time });
}
pub(crate) async fn fail_put_attempt(&self, bucket: &str, object: &str, attempt: usize) {
let mut state = self.state.lock().await;
state
.fail_put_attempts
.entry((bucket.to_string(), object.to_string()))
.or_default()
.insert(attempt);
}
pub(crate) async fn fail_delete_attempt(&self, bucket: &str, object: &str, attempt: usize) {
let mut state = self.state.lock().await;
state
.fail_delete_attempts
.entry((bucket.to_string(), object.to_string()))
.or_default()
.insert(attempt);
}
pub(crate) async fn list_call_count(&self) -> usize {
self.state.lock().await.list_calls
}
pub(crate) async fn read_call_count(&self) -> usize {
self.state.lock().await.read_calls
}
pub(crate) async fn metadata_call_count(&self) -> usize {
self.state.lock().await.metadata_calls
}
pub(crate) async fn reset_call_counts(&self) {
let mut state = self.state.lock().await;
state.read_calls = 0;
state.metadata_calls = 0;
state.list_calls = 0;
}
pub(crate) async fn write_lock_acquisition_count(&self, bucket: &str, object: &str) -> usize {
self.state
.lock()
.await
.write_lock_acquisitions
.get(&(bucket.to_string(), object.to_string()))
.copied()
.unwrap_or_default()
}
pub(crate) async fn read_lock_acquisition_count(&self, bucket: &str, object: &str) -> usize {
self.state
.lock()
.await
.read_lock_acquisitions
.get(&(bucket.to_string(), object.to_string()))
.copied()
.unwrap_or_default()
}
pub(crate) async fn fail_next_read(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.read_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_read_attempts.entry(key).or_default().insert(next_attempt);
}
pub(crate) async fn pause_next_read(&self, bucket: &str, object: &str) -> TestCatalogObjectPause {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.read_attempts.get(&key).copied().unwrap_or_default() + 1;
let pause = TestCatalogObjectPause::default();
state
.pause_read_attempts
.entry(key)
.or_default()
.insert(next_attempt, pause.clone());
pause
}
pub(crate) async fn pause_before_next_read(&self, bucket: &str, object: &str) -> TestCatalogObjectPause {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.read_attempts.get(&key).copied().unwrap_or_default() + 1;
let pause = TestCatalogObjectPause::default();
state
.pause_before_read_attempts
.entry(key)
.or_default()
.insert(next_attempt, pause.clone());
pause
}
pub(crate) async fn omit_etag_for_object(&self, bucket: &str, object: &str) {
self.state
.lock()
.await
.etagless_objects
.insert((bucket.to_string(), object.to_string()));
}
pub(crate) async fn last_read_limit(&self, bucket: &str, object: &str) -> Option<usize> {
let key = (bucket.to_string(), object.to_string());
self.state
.lock()
.await
.read_limits
.iter()
.rev()
.find_map(|(read_key, limit)| (read_key == &key).then_some(*limit))
}
pub(crate) async fn fail_next_put(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_put_attempts.entry(key).or_default().insert(next_attempt);
}
pub(crate) async fn fail_after_next_put(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_after_put_attempts.entry(key).or_default().insert(next_attempt);
}
pub(crate) async fn fail_after_next_delete(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.delete_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_after_delete_attempts.entry(key).or_default().insert(next_attempt);
}
pub(crate) async fn pause_next_put(&self, bucket: &str, object: &str) -> TestCatalogObjectPause {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1;
let pause = TestCatalogObjectPause::default();
state
.pause_put_attempts
.entry(key)
.or_default()
.insert(next_attempt, pause.clone());
pause
}
pub(crate) async fn put_attempt_count(&self, bucket: &str, object: &str) -> usize {
self.state
.lock()
.await
.put_attempts
.get(&(bucket.to_string(), object.to_string()))
.copied()
.unwrap_or_default()
}
}
impl TestCatalogObjectState {
pub(crate) fn next_etag(&mut self) -> String {
self.next_etag += 1;
format!("etag-{}", self.next_etag)
}
}
#[async_trait::async_trait]
impl TableCatalogObjectBackend for TestCatalogObjectBackend {
fn strong_catalog_runtime(&self) -> Option<StrongTableCatalogRuntime> {
self.strong_runtime.clone()
}
async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Option<TableCatalogObject>> {
let mut missing_read_object_path = self.missing_read_object_path.lock().await;
if missing_read_object_path.as_deref() == Some(object) {
missing_read_object_path.take();
return Ok(None);
}
drop(missing_read_object_path);
let mut fail_read_object_path = self.fail_read_object_path.lock().await;
if fail_read_object_path.as_deref() == Some(object) {
fail_read_object_path.take();
return Err(TableCatalogStoreError::Internal("private generated metadata read failure".to_string()));
}
drop(fail_read_object_path);
let key = (bucket.to_string(), object.to_string());
let (attempt, pause_before) = {
let mut state = self.state.lock().await;
state.read_calls += 1;
let attempt = {
let attempts = state.read_attempts.entry(key.clone()).or_default();
*attempts += 1;
*attempts
};
if state
.fail_read_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected read failure for {object} attempt {attempt}"
)));
}
let pause = state
.pause_before_read_attempts
.get_mut(&key)
.and_then(|attempts| attempts.remove(&attempt));
(attempt, pause)
};
if let Some(pause) = pause_before {
pause.started.notify_one();
pause.release.notified().await;
}
let (result, pause) = {
let mut state = self.state.lock().await;
let etagless = state.etagless_objects.contains(&key);
let result = state.objects.get(&key).map(|record| TableCatalogObject {
data: record.data.clone(),
etag: (!etagless).then(|| record.etag.clone()),
mod_time: record.mod_time,
});
let pause = state
.pause_read_attempts
.get_mut(&key)
.and_then(|attempts| attempts.remove(&attempt));
(result, pause)
};
if let Some(pause) = pause {
pause.started.notify_one();
pause.release.notified().await;
}
Ok(result)
}
async fn read_object_limited(
&self,
bucket: &str,
object: &str,
max_size: usize,
) -> TableCatalogStoreResult<Option<TableCatalogObject>> {
self.state
.lock()
.await
.read_limits
.push(((bucket.to_string(), object.to_string()), max_size));
let result = self.read_object(bucket, object).await?;
if result.as_ref().is_some_and(|object| object.data.len() > max_size) {
return Err(TableCatalogStoreError::Invalid(format!(
"catalog object {bucket}/{object} exceeds the maximum size of {max_size} bytes"
)));
}
Ok(result)
}
async fn object_metadata(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Option<TableCatalogObjectMetadata>> {
let mut state = self.state.lock().await;
state.metadata_calls += 1;
let key = (bucket.to_string(), object.to_string());
let etagless = state.etagless_objects.contains(&key);
Ok(state.objects.get(&key).map(|record| TableCatalogObjectMetadata {
etag: (!etagless).then(|| record.etag.clone()),
mod_time: record.mod_time,
}))
}
async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<bool> {
let state = self.state.lock().await;
Ok(state.objects.contains_key(&(bucket.to_string(), object.to_string())))
}
async fn put_object(
&self,
bucket: &str,
object: &str,
data: Vec<u8>,
precondition: TableCatalogPutPrecondition,
) -> TableCatalogStoreResult<()> {
let mut fail_put_object_path = self.fail_put_object_path.lock().await;
if fail_put_object_path.as_deref() == Some(object) {
fail_put_object_path.take();
return Err(TableCatalogStoreError::Internal("injected metadata write failure".to_string()));
}
drop(fail_put_object_path);
let mut corrupt_put_object_path = self.corrupt_put_object_path.lock().await;
let data = if corrupt_put_object_path.as_deref() == Some(object) {
corrupt_put_object_path.take();
b"{}".to_vec()
} else {
data
};
drop(corrupt_put_object_path);
let key = (bucket.to_string(), object.to_string());
let (attempt, pause) = {
let mut state = self.state.lock().await;
let attempt = {
let attempts = state.put_attempts.entry(key.clone()).or_default();
*attempts += 1;
*attempts
};
if state
.fail_put_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected put failure for {object} attempt {attempt}"
)));
}
let pause = state
.pause_put_attempts
.get_mut(&key)
.and_then(|attempts| attempts.remove(&attempt));
(attempt, pause)
};
if let Some(pause) = pause {
pause.started.notify_one();
pause.release.notified().await;
}
let result = {
let mut state = self.state.lock().await;
let precondition_failure = match &precondition {
TableCatalogPutPrecondition::IfAbsent if state.objects.contains_key(&key) => {
Some(TableCatalogStoreError::Conflict(format!("object already exists: {object}")))
}
TableCatalogPutPrecondition::IfMatch(expected) => match state.objects.get(&key) {
None => Some(TableCatalogStoreError::Conflict(format!("object is missing: {object}"))),
Some(current) if &current.etag != expected => {
Some(TableCatalogStoreError::Conflict(format!("object changed: {object}")))
}
Some(_) => None,
},
_ => None,
};
if let Some(err) = precondition_failure {
Err(err)
} else {
let etag = if self.content_addressed_etags {
content_etag(&data)
} else {
state.next_etag()
};
state.objects.insert(
key.clone(),
TestCatalogObjectRecord {
data,
etag,
mod_time: Some(OffsetDateTime::now_utc()),
},
);
if state
.fail_after_put_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
Err(TableCatalogStoreError::Internal(format!(
"injected post-commit put failure for {object} attempt {attempt}"
)))
} else {
Ok(())
}
}
};
if let Some(barrier) = &self.put_object_barrier {
barrier.wait().await;
}
result
}
async fn delete_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<()> {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let attempt = {
let attempts = state.delete_attempts.entry(key.clone()).or_default();
*attempts += 1;
*attempts
};
if state
.fail_delete_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected delete failure for {object} attempt {attempt}"
)));
}
state.objects.remove(&key);
if state
.fail_after_delete_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected post-commit delete failure for {object} attempt {attempt}"
)));
}
Ok(())
}
async fn list_objects(&self, bucket: &str, prefix: &str) -> TableCatalogStoreResult<Vec<String>> {
let mut state = self.state.lock().await;
state.list_calls += 1;
Ok(state
.objects
.keys()
.filter(|(entry_bucket, object)| entry_bucket == bucket && object.starts_with(prefix))
.map(|(_, object)| object.clone())
.collect())
}
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
self.lock_attempts.lock().await.push((bucket.to_string(), object.to_string()));
{
let mut state = self.state.lock().await;
*state
.write_lock_acquisitions
.entry((bucket.to_string(), object.to_string()))
.or_default() += 1;
}
let lock = {
let mut locks = self.locks.lock().await;
locks
.entry((bucket.to_string(), object.to_string()))
.or_insert_with(|| std::sync::Arc::new(tokio::sync::RwLock::new(())))
.clone()
};
Ok(Box::new(lock.write_owned().await))
}
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
// The admin fake implemented only acquire_write_lock, so the trait's
// default read->write delegation made read acquisitions observable in
// lock_attempts as well; keep that (backlog#1837 PR2).
self.lock_attempts.lock().await.push((bucket.to_string(), object.to_string()));
{
let mut state = self.state.lock().await;
*state
.read_lock_acquisitions
.entry((bucket.to_string(), object.to_string()))
.or_default() += 1;
}
let lock = {
let mut locks = self.locks.lock().await;
locks
.entry((bucket.to_string(), object.to_string()))
.or_insert_with(|| std::sync::Arc::new(tokio::sync::RwLock::new(())))
.clone()
};
Ok(Box::new(lock.read_owned().await))
}
}
fn content_etag(data: &[u8]) -> String {
use sha2::Digest;
hex_simd::encode_to_string(sha2::Sha256::digest(data), hex_simd::AsciiCase::Lower)
}
/// Admin-handler-test conveniences carried over from the former
/// TestTableCatalogObjectBackend (backlog#1837 PR2): content-addressed etags
/// (sha256), direct record insertion, and lock observability.
impl TestCatalogObjectBackend {
/// Fake with the admin fixtures' content-addressed etag semantics.
pub(crate) fn content_addressed() -> Self {
Self {
content_addressed_etags: true,
..Self::default()
}
}
pub(crate) async fn put_bytes(&self, bucket: &str, object: &str, data: Vec<u8>) {
let etag = content_etag(&data);
self.state.lock().await.objects.insert(
(bucket.to_string(), object.to_string()),
TestCatalogObjectRecord {
data,
etag,
mod_time: None,
},
);
}
pub(crate) async fn put_json(&self, bucket: &str, object: &str, value: serde_json::Value) {
self.put_json_with_mod_time(bucket, object, value, None).await;
}
pub(crate) async fn put_gzip_json(&self, bucket: &str, object: &str, value: serde_json::Value) {
use std::io::Write;
let data = serde_json::to_vec(&value).expect("metadata JSON should serialize");
let mut encoder = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
encoder.write_all(&data).expect("metadata JSON should compress");
self.put_bytes(bucket, object, encoder.finish().expect("metadata gzip stream should finish"))
.await;
}
pub(crate) async fn put_json_with_mod_time(
&self,
bucket: &str,
object: &str,
value: serde_json::Value,
mod_time: Option<OffsetDateTime>,
) {
let data = serde_json::to_vec(&value).expect("metadata JSON should serialize");
let etag = content_etag(&data);
self.state
.lock()
.await
.objects
.insert((bucket.to_string(), object.to_string()), TestCatalogObjectRecord { data, etag, mod_time });
}
pub(crate) async fn write_lock_is_held(&self, bucket: &str, object: &str) -> bool {
let lock = self
.locks
.lock()
.await
.get(&(bucket.to_string(), object.to_string()))
.cloned();
lock.is_some_and(|lock| lock.try_write_owned().is_err())
}
pub(crate) async fn wait_for_lock_attempts(&self, count: usize) {
tokio::time::timeout(std::time::Duration::from_secs(2), async {
loop {
if self.lock_attempts.lock().await.len() >= count {
return;
}
tokio::task::yield_now().await;
}
})
.await
.expect("lock acquisition attempts should be observable");
}
}
+1 -548
View File
@@ -3,6 +3,7 @@ use super::identifier::{
default_table_lifecycle_path, default_table_marker_path, default_table_root_prefix, is_valid_table_metadata_file_name,
namespace_name_from_marker_path, table_name_from_marker_path, validate_object_mutation,
};
use super::test_support::{BlockingObjectPublication, TestCatalogObjectBackend, UnserializedTestPublication};
use super::*;
use datafusion::{
arrow::{
@@ -1077,312 +1078,6 @@ fn catalog_object_entry_paths_use_internal_root_and_hashed_untrusted_ids() {
}
}
#[derive(Clone, Default)]
struct TestCatalogObjectBackend {
state: Arc<tokio::sync::Mutex<TestCatalogObjectState>>,
locks: TestCatalogObjectLocks,
strong_runtime: Option<StrongTableCatalogRuntime>,
}
type TestCatalogObjectLockKey = (String, String);
type TestCatalogObjectLock = Arc<tokio::sync::RwLock<()>>;
type TestCatalogObjectLocks = Arc<tokio::sync::Mutex<BTreeMap<TestCatalogObjectLockKey, TestCatalogObjectLock>>>;
#[derive(Clone, Default)]
struct TestCatalogObjectPause {
started: Arc<tokio::sync::Notify>,
release: Arc<tokio::sync::Notify>,
}
impl TestCatalogObjectPause {
async fn wait_started(&self) {
self.started.notified().await;
}
fn release(&self) {
self.release.notify_one();
}
}
#[derive(Clone)]
struct BlockingObjectPublication {
backend: TestCatalogObjectBackend,
object: String,
started: Arc<tokio::sync::Notify>,
guard: Arc<parking_lot::Mutex<Option<Box<dyn Send>>>>,
}
impl BlockingObjectPublication {
fn new(backend: TestCatalogObjectBackend, object: impl Into<String>) -> Self {
Self {
backend,
object: object.into(),
started: Arc::new(tokio::sync::Notify::new()),
guard: Arc::new(parking_lot::Mutex::new(None)),
}
}
async fn wait_started(&self) {
self.started.notified().await;
}
}
#[derive(Default)]
struct UnserializedTestPublication;
#[async_trait::async_trait]
impl TableCommitPublication for UnserializedTestPublication {
async fn begin_table_bucket(&self, _table_bucket: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
async fn prepare(&self, _table_bucket: &str, _namespace: &str, _table: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
fn holds_table_bucket(&self, _table_bucket: &str) -> bool {
true
}
fn holds_table(&self, _table_bucket: &str, _namespace: &str, _table: &str) -> bool {
true
}
fn complete(&self) {}
}
#[async_trait::async_trait]
impl TableCommitPublication for BlockingObjectPublication {
async fn begin_table_bucket(&self, _table_bucket: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
async fn prepare(&self, table_bucket: &str, _namespace: &str, _table: &str) -> TableCatalogStoreResult<()> {
self.started.notify_one();
let guard = self.backend.acquire_read_lock(table_bucket, &self.object).await?;
*self.guard.lock() = Some(guard);
Ok(())
}
fn holds_table_bucket(&self, _table_bucket: &str) -> bool {
true
}
fn holds_table(&self, _table_bucket: &str, _namespace: &str, _table: &str) -> bool {
self.guard.lock().is_some()
}
fn complete(&self) {
drop(self.guard.lock().take());
}
}
#[derive(Default)]
struct TestCatalogObjectState {
objects: BTreeMap<(String, String), TestCatalogObjectRecord>,
etagless_objects: BTreeSet<(String, String)>,
fail_read_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pause_before_read_attempts: BTreeMap<(String, String), BTreeMap<usize, TestCatalogObjectPause>>,
pause_read_attempts: BTreeMap<(String, String), BTreeMap<usize, TestCatalogObjectPause>>,
read_attempts: BTreeMap<(String, String), usize>,
read_limits: Vec<((String, String), usize)>,
fail_put_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
fail_after_put_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
pause_put_attempts: BTreeMap<(String, String), BTreeMap<usize, TestCatalogObjectPause>>,
fail_delete_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
fail_after_delete_attempts: BTreeMap<(String, String), BTreeSet<usize>>,
put_attempts: BTreeMap<(String, String), usize>,
delete_attempts: BTreeMap<(String, String), usize>,
write_lock_acquisitions: BTreeMap<(String, String), usize>,
read_lock_acquisitions: BTreeMap<(String, String), usize>,
read_calls: usize,
metadata_calls: usize,
list_calls: usize,
next_etag: u64,
}
#[derive(Clone)]
struct TestCatalogObjectRecord {
data: Vec<u8>,
etag: String,
mod_time: Option<OffsetDateTime>,
}
impl TestCatalogObjectBackend {
async fn seed_object(&self, bucket: &str, object: &str, data: Vec<u8>) {
self.seed_object_with_mod_time(bucket, object, data, Some(OffsetDateTime::UNIX_EPOCH))
.await;
}
async fn seed_object_with_mod_time(&self, bucket: &str, object: &str, data: Vec<u8>, mod_time: Option<OffsetDateTime>) {
let mut state = self.state.lock().await;
let etag = state.next_etag();
state
.objects
.insert((bucket.to_string(), object.to_string()), TestCatalogObjectRecord { data, etag, mod_time });
}
async fn fail_put_attempt(&self, bucket: &str, object: &str, attempt: usize) {
let mut state = self.state.lock().await;
state
.fail_put_attempts
.entry((bucket.to_string(), object.to_string()))
.or_default()
.insert(attempt);
}
async fn fail_delete_attempt(&self, bucket: &str, object: &str, attempt: usize) {
let mut state = self.state.lock().await;
state
.fail_delete_attempts
.entry((bucket.to_string(), object.to_string()))
.or_default()
.insert(attempt);
}
async fn list_call_count(&self) -> usize {
self.state.lock().await.list_calls
}
async fn read_call_count(&self) -> usize {
self.state.lock().await.read_calls
}
async fn metadata_call_count(&self) -> usize {
self.state.lock().await.metadata_calls
}
async fn reset_call_counts(&self) {
let mut state = self.state.lock().await;
state.read_calls = 0;
state.metadata_calls = 0;
state.list_calls = 0;
}
async fn write_lock_acquisition_count(&self, bucket: &str, object: &str) -> usize {
self.state
.lock()
.await
.write_lock_acquisitions
.get(&(bucket.to_string(), object.to_string()))
.copied()
.unwrap_or_default()
}
async fn read_lock_acquisition_count(&self, bucket: &str, object: &str) -> usize {
self.state
.lock()
.await
.read_lock_acquisitions
.get(&(bucket.to_string(), object.to_string()))
.copied()
.unwrap_or_default()
}
async fn fail_next_read(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.read_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_read_attempts.entry(key).or_default().insert(next_attempt);
}
async fn pause_next_read(&self, bucket: &str, object: &str) -> TestCatalogObjectPause {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.read_attempts.get(&key).copied().unwrap_or_default() + 1;
let pause = TestCatalogObjectPause::default();
state
.pause_read_attempts
.entry(key)
.or_default()
.insert(next_attempt, pause.clone());
pause
}
async fn pause_before_next_read(&self, bucket: &str, object: &str) -> TestCatalogObjectPause {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.read_attempts.get(&key).copied().unwrap_or_default() + 1;
let pause = TestCatalogObjectPause::default();
state
.pause_before_read_attempts
.entry(key)
.or_default()
.insert(next_attempt, pause.clone());
pause
}
async fn omit_etag_for_object(&self, bucket: &str, object: &str) {
self.state
.lock()
.await
.etagless_objects
.insert((bucket.to_string(), object.to_string()));
}
async fn last_read_limit(&self, bucket: &str, object: &str) -> Option<usize> {
let key = (bucket.to_string(), object.to_string());
self.state
.lock()
.await
.read_limits
.iter()
.rev()
.find_map(|(read_key, limit)| (read_key == &key).then_some(*limit))
}
async fn fail_next_put(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_put_attempts.entry(key).or_default().insert(next_attempt);
}
async fn fail_after_next_put(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_after_put_attempts.entry(key).or_default().insert(next_attempt);
}
async fn fail_after_next_delete(&self, bucket: &str, object: &str) {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.delete_attempts.get(&key).copied().unwrap_or_default() + 1;
state.fail_after_delete_attempts.entry(key).or_default().insert(next_attempt);
}
async fn pause_next_put(&self, bucket: &str, object: &str) -> TestCatalogObjectPause {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1;
let pause = TestCatalogObjectPause::default();
state
.pause_put_attempts
.entry(key)
.or_default()
.insert(next_attempt, pause.clone());
pause
}
async fn put_attempt_count(&self, bucket: &str, object: &str) -> usize {
self.state
.lock()
.await
.put_attempts
.get(&(bucket.to_string(), object.to_string()))
.copied()
.unwrap_or_default()
}
}
impl TestCatalogObjectState {
fn next_etag(&mut self) -> String {
self.next_etag += 1;
format!("etag-{}", self.next_etag)
}
}
fn maintenance_object_report<'a>(
report: &'a TableMetadataMaintenanceReport,
metadata_location: &str,
@@ -2466,248 +2161,6 @@ fn parquet_i32_values(data: Vec<u8>) -> Vec<i32> {
values
}
#[async_trait::async_trait]
impl TableCatalogObjectBackend for TestCatalogObjectBackend {
fn strong_catalog_runtime(&self) -> Option<StrongTableCatalogRuntime> {
self.strong_runtime.clone()
}
async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Option<TableCatalogObject>> {
let key = (bucket.to_string(), object.to_string());
let (attempt, pause_before) = {
let mut state = self.state.lock().await;
state.read_calls += 1;
let attempt = {
let attempts = state.read_attempts.entry(key.clone()).or_default();
*attempts += 1;
*attempts
};
if state
.fail_read_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected read failure for {object} attempt {attempt}"
)));
}
let pause = state
.pause_before_read_attempts
.get_mut(&key)
.and_then(|attempts| attempts.remove(&attempt));
(attempt, pause)
};
if let Some(pause) = pause_before {
pause.started.notify_one();
pause.release.notified().await;
}
let (result, pause) = {
let mut state = self.state.lock().await;
let etagless = state.etagless_objects.contains(&key);
let result = state.objects.get(&key).map(|record| TableCatalogObject {
data: record.data.clone(),
etag: (!etagless).then(|| record.etag.clone()),
mod_time: record.mod_time,
});
let pause = state
.pause_read_attempts
.get_mut(&key)
.and_then(|attempts| attempts.remove(&attempt));
(result, pause)
};
if let Some(pause) = pause {
pause.started.notify_one();
pause.release.notified().await;
}
Ok(result)
}
async fn read_object_limited(
&self,
bucket: &str,
object: &str,
max_size: usize,
) -> TableCatalogStoreResult<Option<TableCatalogObject>> {
self.state
.lock()
.await
.read_limits
.push(((bucket.to_string(), object.to_string()), max_size));
let result = self.read_object(bucket, object).await?;
if result.as_ref().is_some_and(|object| object.data.len() > max_size) {
return Err(TableCatalogStoreError::Invalid(format!(
"catalog object {bucket}/{object} exceeds the maximum size of {max_size} bytes"
)));
}
Ok(result)
}
async fn object_metadata(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Option<TableCatalogObjectMetadata>> {
let mut state = self.state.lock().await;
state.metadata_calls += 1;
let key = (bucket.to_string(), object.to_string());
let etagless = state.etagless_objects.contains(&key);
Ok(state.objects.get(&key).map(|record| TableCatalogObjectMetadata {
etag: (!etagless).then(|| record.etag.clone()),
mod_time: record.mod_time,
}))
}
async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<bool> {
let state = self.state.lock().await;
Ok(state.objects.contains_key(&(bucket.to_string(), object.to_string())))
}
async fn put_object(
&self,
bucket: &str,
object: &str,
data: Vec<u8>,
precondition: TableCatalogPutPrecondition,
) -> TableCatalogStoreResult<()> {
let key = (bucket.to_string(), object.to_string());
let (attempt, pause) = {
let mut state = self.state.lock().await;
let attempt = {
let attempts = state.put_attempts.entry(key.clone()).or_default();
*attempts += 1;
*attempts
};
if state
.fail_put_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected put failure for {object} attempt {attempt}"
)));
}
let pause = state
.pause_put_attempts
.get_mut(&key)
.and_then(|attempts| attempts.remove(&attempt));
(attempt, pause)
};
if let Some(pause) = pause {
pause.started.notify_one();
pause.release.notified().await;
}
let mut state = self.state.lock().await;
match precondition {
TableCatalogPutPrecondition::IfAbsent if state.objects.contains_key(&key) => {
return Err(TableCatalogStoreError::Conflict(format!("object already exists: {object}")));
}
TableCatalogPutPrecondition::IfMatch(expected) => {
let Some(current) = state.objects.get(&key) else {
return Err(TableCatalogStoreError::Conflict(format!("object is missing: {object}")));
};
if current.etag != expected {
return Err(TableCatalogStoreError::Conflict(format!("object changed: {object}")));
}
}
_ => {}
}
let etag = state.next_etag();
state.objects.insert(
key.clone(),
TestCatalogObjectRecord {
data,
etag,
mod_time: Some(OffsetDateTime::now_utc()),
},
);
if state
.fail_after_put_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected post-commit put failure for {object} attempt {attempt}"
)));
}
Ok(())
}
async fn delete_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<()> {
let mut state = self.state.lock().await;
let key = (bucket.to_string(), object.to_string());
let attempt = {
let attempts = state.delete_attempts.entry(key.clone()).or_default();
*attempts += 1;
*attempts
};
if state
.fail_delete_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected delete failure for {object} attempt {attempt}"
)));
}
state.objects.remove(&key);
if state
.fail_after_delete_attempts
.get(&key)
.is_some_and(|attempts| attempts.contains(&attempt))
{
return Err(TableCatalogStoreError::Internal(format!(
"injected post-commit delete failure for {object} attempt {attempt}"
)));
}
Ok(())
}
async fn list_objects(&self, bucket: &str, prefix: &str) -> TableCatalogStoreResult<Vec<String>> {
let mut state = self.state.lock().await;
state.list_calls += 1;
Ok(state
.objects
.keys()
.filter(|(entry_bucket, object)| entry_bucket == bucket && object.starts_with(prefix))
.map(|(_, object)| object.clone())
.collect())
}
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
{
let mut state = self.state.lock().await;
*state
.write_lock_acquisitions
.entry((bucket.to_string(), object.to_string()))
.or_default() += 1;
}
let lock = {
let mut locks = self.locks.lock().await;
locks
.entry((bucket.to_string(), object.to_string()))
.or_insert_with(|| std::sync::Arc::new(tokio::sync::RwLock::new(())))
.clone()
};
Ok(Box::new(lock.write_owned().await))
}
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
{
let mut state = self.state.lock().await;
*state
.read_lock_acquisitions
.entry((bucket.to_string(), object.to_string()))
.or_default() += 1;
}
let lock = {
let mut locks = self.locks.lock().await;
locks
.entry((bucket.to_string(), object.to_string()))
.or_insert_with(|| std::sync::Arc::new(tokio::sync::RwLock::new(())))
.clone()
};
Ok(Box::new(lock.read_owned().await))
}
}
fn test_bucket_entry(bucket: &str) -> TableBucketEntry {
TableBucketEntry {
version: TABLE_CATALOG_ENTRY_VERSION,