diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 88125cbac..853eff9fd 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -140,6 +140,21 @@ struct CoalescedReadVersionRequest { tx: oneshot::Sender>, } +#[derive(Clone, Debug, Eq, PartialEq)] +struct ExpectedBatchReadVersionItem { + path: String, + version_id: String, +} + +impl From<&BatchReadVersionItem> for ExpectedBatchReadVersionItem { + fn from(item: &BatchReadVersionItem) -> Self { + Self { + path: item.path.clone(), + version_id: item.version_id.clone(), + } + } +} + #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] struct ReadVersionCoalescerKey { disk: usize, @@ -266,7 +281,7 @@ async fn flush_read_version_coalescer_pending( items.push(request.item); } - let expected_items = items.clone(); + let expected_items = items.iter().map(ExpectedBatchReadVersionItem::from).collect::>(); record_read_version_coalescer_event("attempted_batch", items.len()); let result = match tokio::time::timeout(get_drive_metadata_timeout(), disk.batch_read_version(BatchReadVersionReq { items, opts })) @@ -292,7 +307,7 @@ async fn flush_read_version_coalescer_pending( } fn map_batch_read_version_responses( - expected_items: &[BatchReadVersionItem], + expected_items: &[ExpectedBatchReadVersionItem], responses: Vec, ) -> Vec> { let mut results = (0..expected_items.len()) @@ -6878,6 +6893,7 @@ mod tests { }, ]; + let expected_items = expected_batch_read_version_items(&expected_items); let mut results = map_batch_read_version_responses(&expected_items, responses).into_iter(); let first = results .next() @@ -6918,6 +6934,7 @@ mod tests { version_id: "v-b".to_string(), }, ]; + let expected_items = expected_batch_read_version_items(&expected_items); let results = map_batch_read_version_responses( &expected_items, vec![ @@ -6957,6 +6974,7 @@ mod tests { path: "object-a".to_string(), version_id: "v-a".to_string(), }]; + let expected_items = expected_batch_read_version_items(&expected_items); let mismatched = map_batch_read_version_responses( &expected_items, vec![BatchReadVersionResp { @@ -7018,6 +7036,10 @@ mod tests { ); } + fn expected_batch_read_version_items(items: &[BatchReadVersionItem]) -> Vec { + items.iter().map(ExpectedBatchReadVersionItem::from).collect() + } + /// Isolation guard: unobserved objects record nothing (so parallel tests do /// not inflate one another), and a scope clears its own counts on drop. #[tokio::test] diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 7a8653415..e44e8dedf 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -1203,7 +1203,8 @@ fn inject_batch_delete_pool_errors( bucket: &str, pool_idx: usize, object_names: &[String], - result: &mut (Vec, Vec>), + deleted: &[DeletedObject], + errors: &mut [Option], ) { let state = BATCH_DELETE_POOL_ERROR_INJECTION .get_or_init(|| std::sync::Mutex::new(None)) @@ -1220,8 +1221,8 @@ fn inject_batch_delete_pool_errors( let Some(error) = state.errors.get(object_name) else { continue; }; - if result.1[idx].is_none() && result.0[idx].found { - result.1[idx] = Some(error.clone()); + if errors[idx].is_none() && deleted[idx].found { + errors[idx] = Some(error.clone()); state.observed.fetch_add(1, Ordering::AcqRel); } } @@ -3194,7 +3195,7 @@ impl ECStore { // Default return value let mut del_objects = vec![DeletedObject::default(); objects.len()]; - let accounting = vec![None; objects.len()]; + let mut accounting = vec![None; objects.len()]; let mut del_errs = Vec::with_capacity(objects.len()); for _ in 0..objects.len() { @@ -3333,11 +3334,12 @@ impl ECStore { .iter() .map(|object| object.object_name.clone()) .collect::>(); - let result = pool.delete_objects(bucket, pool_objects, pool_opts).await; + let result = pool.delete_objects_with_accounting(bucket, pool_objects, pool_opts).await; #[cfg(test)] let result = { let mut result = result; - inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, &mut result); + let (deleted, errors, _) = &mut result; + inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, deleted, errors); result }; (object_indices, result) @@ -3347,7 +3349,7 @@ impl ECStore { let results = join_all(futures).await; for idx in 0..del_objects.len() { - let pool_results = results.iter().filter_map(|(object_indices, (dels, errs))| { + let pool_results = results.iter().filter_map(|(object_indices, (dels, errs, _))| { let pool_object_idx = object_indices.binary_search(&idx).ok()?; Some((&dels[pool_object_idx], &errs[pool_object_idx])) }); @@ -3367,6 +3369,12 @@ impl ECStore { } } + for (object_indices, (_, _, pool_accounting)) in &results { + for (pool_object_idx, object_idx) in object_indices.iter().enumerate() { + accounting[*object_idx] = pool_accounting.get(pool_object_idx).cloned().flatten(); + } + } + #[cfg(test)] for (idx, object) in objects.iter().enumerate() { if del_errs[idx].is_none() && del_objects[idx].delete_marker {