mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 82bff87f19 | |||
| 78b6879043 | |||
| 520b93c5fc |
@@ -5143,9 +5143,6 @@ impl ECStore {
|
||||
let buckets = self.get_buckets_to_decommission().await?;
|
||||
let pool = self.pools[idx].clone();
|
||||
|
||||
self.ensure_decommission_multipart_uploads_drained(idx, &pool, &buckets)
|
||||
.await?;
|
||||
|
||||
for (set_index, set) in pool.disk_set.iter().enumerate() {
|
||||
for bucket_info in &buckets {
|
||||
let mut lifecycle_config = None;
|
||||
@@ -5259,49 +5256,6 @@ impl ECStore {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn ensure_decommission_multipart_uploads_drained(
|
||||
&self,
|
||||
idx: usize,
|
||||
pool: &Sets,
|
||||
buckets: &[DecomBucketInfo],
|
||||
) -> Result<()> {
|
||||
let mut bucket_names = buckets
|
||||
.iter()
|
||||
.filter(|bucket| bucket.name != RUSTFS_META_BUCKET)
|
||||
.map(|bucket| bucket.name.as_str())
|
||||
.collect::<Vec<_>>();
|
||||
bucket_names.sort_unstable();
|
||||
bucket_names.dedup();
|
||||
|
||||
// Take one bucket fence at a time so cross-bucket COPY cannot form an
|
||||
// ABBA cycle. Suspension prevents new source uploads after each fence.
|
||||
for bucket in bucket_names {
|
||||
let lifecycle_guard = self.acquire_bucket_lifecycle_write_lock(bucket).await?;
|
||||
if lifecycle_guard.is_lock_lost() {
|
||||
return Err(Error::other(format!(
|
||||
"decommission multipart drain lost the bucket lifecycle fence for `{bucket}`"
|
||||
)));
|
||||
}
|
||||
for set in &pool.disk_set {
|
||||
if let Some(upload_path) = set.first_multipart_upload_path_for_decommission(bucket).await? {
|
||||
return Err(Error::other(format!(
|
||||
"pool {idx} still contains multipart upload `{upload_path}` for bucket `{bucket}`; resolve it before retrying decommission"
|
||||
)));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn ensure_decommission_multipart_uploads_drained_for_test(self: &Arc<Self>, idx: usize) -> Result<()> {
|
||||
let buckets = self.get_buckets_to_decommission().await?;
|
||||
let pool = self.pools[idx].clone();
|
||||
self.ensure_decommission_multipart_uploads_drained(idx, pool.as_ref(), &buckets)
|
||||
.await
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self, rd))]
|
||||
async fn decommission_object(
|
||||
self: Arc<Self>,
|
||||
|
||||
@@ -556,69 +556,6 @@ async fn multipart_upload_paths_on_disk(disk: DiskStore, bucket: &str) -> disk::
|
||||
}
|
||||
|
||||
impl SetDisks {
|
||||
async fn discover_multipart_upload_paths(
|
||||
&self,
|
||||
orig_bucket: &str,
|
||||
error_path: &str,
|
||||
) -> Result<(Vec<Option<DiskStore>>, Vec<String>, usize)> {
|
||||
let disks = self.disks.read().await.clone();
|
||||
if disks.is_empty() {
|
||||
return Err(Error::ErasureReadQuorum);
|
||||
}
|
||||
let discovery_quorum = if self.default_parity_count == 0 {
|
||||
disks.len()
|
||||
} else {
|
||||
(disks.len() / 2).max(1)
|
||||
};
|
||||
let mut discovery_errors = (0..disks.len()).map(|_| Some(DiskError::DiskNotFound)).collect::<Vec<_>>();
|
||||
let mut candidate_counts = HashMap::<String, usize>::new();
|
||||
let mut discovery_tasks = JoinSet::new();
|
||||
for (index, disk) in disks.iter().enumerate() {
|
||||
let disk = disk.clone();
|
||||
let orig_bucket = orig_bucket.to_string();
|
||||
discovery_tasks.spawn(async move {
|
||||
let result = match disk {
|
||||
Some(disk) => multipart_upload_paths_on_disk(disk, &orig_bucket).await,
|
||||
None => Err(DiskError::DiskNotFound),
|
||||
};
|
||||
(index, result)
|
||||
});
|
||||
}
|
||||
|
||||
while let Some(task_result) = discovery_tasks.join_next().await {
|
||||
let Ok((index, result)) = task_result else {
|
||||
continue;
|
||||
};
|
||||
match result {
|
||||
Ok(paths) => {
|
||||
discovery_errors[index] = None;
|
||||
for path in paths {
|
||||
*candidate_counts.entry(path).or_insert(0) += 1;
|
||||
}
|
||||
}
|
||||
Err(err) => discovery_errors[index] = Some(err),
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(err) = reduce_read_quorum_errs(&discovery_errors, OBJECT_OP_IGNORED_ERRS, discovery_quorum) {
|
||||
return Err(to_object_err(err.into(), vec![orig_bucket, error_path]));
|
||||
}
|
||||
|
||||
let mut candidate_paths = candidate_counts
|
||||
.into_iter()
|
||||
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
|
||||
.collect::<Vec<_>>();
|
||||
candidate_paths.sort_unstable();
|
||||
Ok((disks, candidate_paths, discovery_quorum))
|
||||
}
|
||||
|
||||
pub(crate) async fn first_multipart_upload_path_for_decommission(&self, bucket: &str) -> Result<Option<String>> {
|
||||
let (_, paths, _) = self
|
||||
.discover_multipart_upload_paths(bucket, RUSTFS_META_MULTIPART_BUCKET)
|
||||
.await?;
|
||||
Ok(paths.into_iter().next())
|
||||
}
|
||||
|
||||
async fn acquire_multipart_upload_read_lock(
|
||||
&self,
|
||||
op: &'static str,
|
||||
@@ -810,7 +747,53 @@ impl SetDisks {
|
||||
max_uploads: usize,
|
||||
expected_incarnation_id: Option<Uuid>,
|
||||
) -> Result<ListMultipartsInfo> {
|
||||
let (disks, candidate_paths, discovery_quorum) = self.discover_multipart_upload_paths(bucket, prefix).await?;
|
||||
let disks = self.disks.read().await.clone();
|
||||
if disks.is_empty() {
|
||||
return Err(Error::ErasureReadQuorum);
|
||||
}
|
||||
let discovery_quorum = if self.default_parity_count == 0 {
|
||||
disks.len()
|
||||
} else {
|
||||
(disks.len() / 2).max(1)
|
||||
};
|
||||
let mut discovery_errors = (0..disks.len()).map(|_| Some(DiskError::DiskNotFound)).collect::<Vec<_>>();
|
||||
let mut candidate_counts = HashMap::<String, usize>::new();
|
||||
let mut discovery_tasks = JoinSet::new();
|
||||
for (index, disk) in disks.iter().enumerate() {
|
||||
let disk = disk.clone();
|
||||
let bucket = bucket.to_string();
|
||||
discovery_tasks.spawn(async move {
|
||||
let result = match disk {
|
||||
Some(disk) => multipart_upload_paths_on_disk(disk, &bucket).await,
|
||||
None => Err(DiskError::DiskNotFound),
|
||||
};
|
||||
(index, result)
|
||||
});
|
||||
}
|
||||
|
||||
while let Some(task_result) = discovery_tasks.join_next().await {
|
||||
let Ok((index, result)) = task_result else {
|
||||
continue;
|
||||
};
|
||||
match result {
|
||||
Ok(paths) => {
|
||||
discovery_errors[index] = None;
|
||||
for path in paths {
|
||||
*candidate_counts.entry(path).or_insert(0) += 1;
|
||||
}
|
||||
}
|
||||
Err(err) => discovery_errors[index] = Some(err),
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(err) = reduce_read_quorum_errs(&discovery_errors, OBJECT_OP_IGNORED_ERRS, discovery_quorum) {
|
||||
return Err(to_object_err(err.into(), vec![bucket, prefix]));
|
||||
}
|
||||
|
||||
let candidate_paths = candidate_counts
|
||||
.into_iter()
|
||||
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
|
||||
.collect::<Vec<_>>();
|
||||
let listed_uploads = stream::iter(candidate_paths)
|
||||
.map(|upload_path| {
|
||||
let disks = &disks;
|
||||
|
||||
@@ -1589,177 +1589,6 @@ mod tests {
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn suspended_decommission_source_multipart_remains_operable_until_drained() {
|
||||
let temp_dir = tempfile::tempdir().expect("create decommission multipart drain store dir");
|
||||
let (_ctx, store, shutdown) =
|
||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-multipart-drain", &[4, 4])).await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let bucket = format!("decommission-multipart-drain-{}", uuid::Uuid::new_v4());
|
||||
let complete_object = "complete.bin";
|
||||
let abort_object = "abort.bin";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create decommission multipart drain bucket");
|
||||
|
||||
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
|
||||
let lifecycle_guard = store
|
||||
.acquire_bucket_lifecycle_read_lock(&bucket)
|
||||
.await
|
||||
.expect("acquire multipart creation lifecycle fence");
|
||||
let mut upload_opts = ObjectOptions {
|
||||
expected_bucket_incarnation_id: Some(incarnation),
|
||||
..Default::default()
|
||||
};
|
||||
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
|
||||
let complete_upload = store.pools[0]
|
||||
.new_multipart_upload(&bucket, complete_object, &upload_opts)
|
||||
.await
|
||||
.expect("create source upload to complete");
|
||||
let abort_upload = store.pools[0]
|
||||
.new_multipart_upload(&bucket, abort_object, &upload_opts)
|
||||
.await
|
||||
.expect("create source upload to abort");
|
||||
drop(lifecycle_guard);
|
||||
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
|
||||
let err = store
|
||||
.ensure_decommission_multipart_uploads_drained_for_test(0)
|
||||
.await
|
||||
.expect_err("an unresolved source multipart upload must block final decommission");
|
||||
let drain_error = err.to_string();
|
||||
assert!(
|
||||
drain_error.contains("still contains multipart upload") && drain_error.contains(&bucket),
|
||||
"the drain error must identify both the upload path and user bucket: {drain_error}"
|
||||
);
|
||||
|
||||
let listed = store
|
||||
.list_multipart_uploads(&bucket, "", None, None, None, 100)
|
||||
.await
|
||||
.expect("list uploads from suspended decommission source");
|
||||
assert!(
|
||||
listed
|
||||
.uploads
|
||||
.iter()
|
||||
.any(|upload| upload.upload_id.as_str() == complete_upload.upload_id.as_str()),
|
||||
"the upload selected before suspension must remain visible"
|
||||
);
|
||||
store
|
||||
.get_multipart_info(&bucket, complete_object, &complete_upload.upload_id, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("read upload metadata from suspended decommission source");
|
||||
|
||||
let mut part_reader = PutObjReader::from_vec(b"multipart body".to_vec());
|
||||
let part = store
|
||||
.put_object_part(
|
||||
&bucket,
|
||||
complete_object,
|
||||
&complete_upload.upload_id,
|
||||
1,
|
||||
&mut part_reader,
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
.expect("write part to suspended decommission source");
|
||||
let parts = store
|
||||
.list_object_parts(&bucket, complete_object, &complete_upload.upload_id, None, 100, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("list parts from suspended decommission source");
|
||||
assert_eq!(parts.parts.len(), 1);
|
||||
assert_eq!(parts.parts[0].etag.as_deref(), part.etag.as_deref());
|
||||
|
||||
store
|
||||
.clone()
|
||||
.complete_multipart_upload(
|
||||
&bucket,
|
||||
complete_object,
|
||||
&complete_upload.upload_id,
|
||||
vec![crate::storage_api_contracts::multipart::CompletePart {
|
||||
part_num: part.part_num,
|
||||
etag: part.etag,
|
||||
..Default::default()
|
||||
}],
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
.expect("complete upload on suspended decommission source");
|
||||
store
|
||||
.abort_multipart_upload(&bucket, abort_object, &abort_upload.upload_id, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("abort upload on suspended decommission source");
|
||||
|
||||
store
|
||||
.ensure_decommission_multipart_uploads_drained_for_test(0)
|
||||
.await
|
||||
.expect("final decommission gate should open after all source uploads are resolved");
|
||||
assert_pool_object_present(&store.pools[0], &bucket, complete_object).await;
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn active_multipart_upload_routes_before_faulted_suspended_source() {
|
||||
let temp_dir = tempfile::tempdir().expect("create active-first multipart routing store dir");
|
||||
let (_ctx, store, shutdown) =
|
||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "active-first-multipart-routing", &[4, 4]))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let bucket = format!("active-first-multipart-routing-{}", uuid::Uuid::new_v4());
|
||||
let object = "target-upload.bin";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create active-first multipart routing bucket");
|
||||
|
||||
let incarnation = store.bucket_incarnation_id(&bucket).await.expect("read bucket incarnation");
|
||||
let lifecycle_guard = store
|
||||
.acquire_bucket_lifecycle_read_lock(&bucket)
|
||||
.await
|
||||
.expect("acquire multipart creation lifecycle fence");
|
||||
let mut upload_opts = ObjectOptions {
|
||||
expected_bucket_incarnation_id: Some(incarnation),
|
||||
..Default::default()
|
||||
};
|
||||
upload_opts.add_bucket_lifecycle_lock_guard(&lifecycle_guard);
|
||||
let upload = store.pools[1]
|
||||
.new_multipart_upload(&bucket, object, &upload_opts)
|
||||
.await
|
||||
.expect("create upload in active target pool");
|
||||
drop(lifecycle_guard);
|
||||
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
let source_set = store.pools[0].get_disks_by_key(object);
|
||||
let original_source_disks = {
|
||||
let mut disks = source_set.disks.write().await;
|
||||
let original = disks.clone();
|
||||
disks.fill(None);
|
||||
original
|
||||
};
|
||||
|
||||
let source_result = store.pools[0]
|
||||
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
|
||||
.await;
|
||||
let routed_result = store
|
||||
.get_multipart_info(&bucket, object, &upload.upload_id, &ObjectOptions::default())
|
||||
.await;
|
||||
*source_set.disks.write().await = original_source_disks;
|
||||
|
||||
assert!(
|
||||
matches!(&source_result, Err(StorageError::ErasureReadQuorum)),
|
||||
"the suspended source must expose the injected hard read failure: {source_result:?}"
|
||||
);
|
||||
let routed = routed_result.expect("the active target UploadID must be resolved before the faulted suspended source");
|
||||
assert_eq!(routed.upload_id, upload.upload_id);
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn delete_objects_skips_active_rebalance_source_pool() {
|
||||
|
||||
@@ -196,25 +196,6 @@ async fn list_pool_multipart_uploads_for_incarnation(
|
||||
}
|
||||
|
||||
impl ECStore {
|
||||
async fn existing_multipart_pool_order(&self) -> Vec<usize> {
|
||||
// A draining source must not hide a valid UploadID in an active target,
|
||||
// while physical order within each phase preserves fail-closed errors.
|
||||
let mut active = Vec::with_capacity(self.pools.len());
|
||||
let mut draining = Vec::new();
|
||||
for (idx, pool) in self.pools.iter().enumerate() {
|
||||
if self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
if self.is_suspended(pool.pool_idx).await {
|
||||
draining.push(idx);
|
||||
} else {
|
||||
active.push(idx);
|
||||
}
|
||||
}
|
||||
active.extend(draining);
|
||||
active
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn list_multipart_uploads_for_bucket_incarnation(
|
||||
&self,
|
||||
@@ -309,8 +290,10 @@ impl ECStore {
|
||||
.await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
return match pool
|
||||
.list_object_parts(bucket, object, upload_id, part_number_marker, max_parts, opts)
|
||||
.await
|
||||
@@ -370,8 +353,10 @@ impl ECStore {
|
||||
let mut common_prefixes = HashSet::new();
|
||||
let mut source_truncated = false;
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
let res = list_pool_multipart_uploads_for_incarnation(
|
||||
pool,
|
||||
bucket,
|
||||
@@ -538,8 +523,10 @@ impl ECStore {
|
||||
.await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await {
|
||||
Ok(res) => return Ok(res),
|
||||
Err(err) => {
|
||||
@@ -599,8 +586,10 @@ impl ECStore {
|
||||
return self.pools[0].get_multipart_info(bucket, object, upload_id, opts).await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
return match pool.get_multipart_info(bucket, object, upload_id, opts).await {
|
||||
Ok(res) => Ok(res),
|
||||
@@ -635,8 +624,10 @@ impl ECStore {
|
||||
return self.pools[0].abort_multipart_upload(bucket, object, upload_id, opts).await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
let err = match pool.abort_multipart_upload(bucket, object, upload_id, opts).await {
|
||||
Ok(_) => return Ok(()),
|
||||
@@ -694,8 +685,10 @@ impl ECStore {
|
||||
.await;
|
||||
}
|
||||
|
||||
for pool_idx in self.existing_multipart_pool_order().await {
|
||||
let pool = &self.pools[pool_idx];
|
||||
for pool in self.pools.iter() {
|
||||
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
|
||||
continue;
|
||||
}
|
||||
|
||||
let pool = pool.clone();
|
||||
let err = match pool
|
||||
|
||||
@@ -34,8 +34,8 @@ CONCURRENCY=8
|
||||
DURATION="60s"
|
||||
ROUNDS=3
|
||||
COOLDOWN_SECS=20
|
||||
DATASET_SETUP_DURATION="10s"
|
||||
HEALTH_TIMEOUT_SECS=180
|
||||
DATASET_OBJECTS_PER_WORKER=8
|
||||
FAIL_PCT=10
|
||||
WARN_PCT=5
|
||||
ALLOW_REGRESSION=false
|
||||
@@ -100,9 +100,6 @@ Benchmark:
|
||||
--duration <dur> warp duration per cell (default 60s).
|
||||
--rounds <n> rounds per cell; must be >= 3 (default 3).
|
||||
--cooldown <n> cooldown seconds between rounds/sizes (default 20).
|
||||
--dataset-setup-duration <dur>
|
||||
isolated Warp PUT warm-up for get/mixed legs
|
||||
(default 10s; not included in the measurement).
|
||||
--concurrency <n> warp concurrency (default 8).
|
||||
--warp-bin <path> warp binary (default warp).
|
||||
|
||||
@@ -173,7 +170,6 @@ while [[ $# -gt 0 ]]; do
|
||||
--duration) DURATION="$2"; shift 2 ;;
|
||||
--rounds) ROUNDS="$2"; shift 2 ;;
|
||||
--cooldown) COOLDOWN_SECS="$2"; shift 2 ;;
|
||||
--dataset-setup-duration) DATASET_SETUP_DURATION="$2"; shift 2 ;;
|
||||
--health-timeout) HEALTH_TIMEOUT_SECS="$2"; shift 2 ;;
|
||||
--fail-pct) FAIL_PCT="$2"; shift 2 ;;
|
||||
--warn-pct) WARN_PCT="$2"; shift 2 ;;
|
||||
@@ -366,29 +362,18 @@ measure() {
|
||||
--duration "$DURATION" --rounds "$ROUNDS" --cooldown-secs "$COOLDOWN_SECS"
|
||||
--out-dir "$cell"
|
||||
)
|
||||
[[ "$mode" == "put" ]] || args+=(--extra-args "--noclear")
|
||||
if [[ "$mode" != "put" ]]; then
|
||||
# Warp defaults to 2,500 setup objects per round. At 10 MiB that writes
|
||||
# 25 GiB before every 12-second measurement, so the matrix cannot finish
|
||||
# inside the workflow budget. Eight objects per worker keeps preparation
|
||||
# bounded while retaining a multi-object working set for relative A/B.
|
||||
args+=(--extra-args "--objects $((CONCURRENCY * DATASET_OBJECTS_PER_WORKER)) --noclear")
|
||||
fi
|
||||
[[ -n "$baseline_csv" ]] && args+=(--baseline-csv "$baseline_csv")
|
||||
run "$ENHANCED_BENCH" "${args[@]}" >&2
|
||||
echo "$cell"
|
||||
}
|
||||
|
||||
prepare_dataset() {
|
||||
local leg="$1" workload="$2" mode="$3" size="$4" sync_label="$5" bucket="$6"
|
||||
[[ "$mode" != "put" ]] || return 0
|
||||
|
||||
local setup_cell="$OUT_DIR/$workload/$sync_label/$leg/dataset-setup"
|
||||
local args=(
|
||||
--tool warp --warp-bin "$WARP_BIN" --warp-mode put
|
||||
--endpoint "$ADDRESS" --access-key "$ACCESS_KEY" --secret-key "$SECRET_KEY"
|
||||
--region "$REGION" --bucket "$bucket" --sizes "$size" --concurrency "$CONCURRENCY"
|
||||
--duration "$DATASET_SETUP_DURATION" --rounds 1 --cooldown-secs 0
|
||||
--extra-args "--noclear"
|
||||
--out-dir "$setup_cell"
|
||||
)
|
||||
log "preparing isolated dataset: $sync_label/$workload/$leg bucket=$bucket"
|
||||
run "$ENHANCED_BENCH" "${args[@]}" >&2
|
||||
}
|
||||
|
||||
write_schedule_header() {
|
||||
echo "sync_label,drive_sync,workload,mode,size,leg,phase,binary,out_dir,bucket,dataset_setup" >"$OUT_DIR/abba_schedule.csv"
|
||||
}
|
||||
@@ -399,7 +384,7 @@ append_schedule() {
|
||||
phase="$(phase_for_leg "$leg")"
|
||||
bin="$(binary_for_leg "$leg")"
|
||||
local dataset_setup="none"
|
||||
[[ "$mode" == "put" ]] || dataset_setup="warp-put"
|
||||
[[ "$mode" == "put" ]] || dataset_setup="warp-native-bounded"
|
||||
echo "$sync_label,$drive_sync,$workload,$mode,$size,$leg,$phase,$bin,$OUT_DIR/$workload/$sync_label/$leg,$bucket,$dataset_setup" >>"$OUT_DIR/abba_schedule.csv"
|
||||
}
|
||||
|
||||
@@ -481,7 +466,8 @@ dataset_namespace=$DATASET_NAMESPACE
|
||||
local_run_data_root=$RUN_DATA_ROOT
|
||||
bucket_isolation=per-leg
|
||||
bucket_prefix=rustfs-abba-$DATASET_NAMESPACE
|
||||
dataset_setup=get-and-mixed-via-warp-put
|
||||
dataset_setup=get-and-mixed-via-bounded-warp-native
|
||||
dataset_objects=$((CONCURRENCY * DATASET_OBJECTS_PER_WORKER))
|
||||
endpoint=$ADDRESS
|
||||
warp_version=$("$WARP_BIN" --version 2>/dev/null | head -n1 || echo unknown)
|
||||
EOF
|
||||
@@ -502,7 +488,6 @@ for ds_spec in "${DRIVE_SYNC_MATRIX[@]}"; do
|
||||
log "=== $sync_label $workload leg $leg ($(phase_for_leg "$leg")) ==="
|
||||
bucket="$(bucket_for_leg "$sync_label" "$workload" "$leg")"
|
||||
bring_up "$leg" "$drive_sync" "$workload" "$mode" "$size" "$sync_label" "$bucket"
|
||||
prepare_dataset "$leg" "$workload" "$mode" "$size" "$sync_label" "$bucket"
|
||||
append_schedule "$sync_label" "$drive_sync" "$workload" "$mode" "$size" "$leg" "$bucket"
|
||||
|
||||
baseline_csv=""
|
||||
|
||||
@@ -672,7 +672,7 @@ extract_report_line() {
|
||||
local regex="$1"
|
||||
local file="$2"
|
||||
awk -v regex="$regex" '
|
||||
/^Report:/ {
|
||||
/^(Report|Operation):/ {
|
||||
in_report = 1
|
||||
next
|
||||
}
|
||||
@@ -703,9 +703,9 @@ normalize_duration_metric() {
|
||||
extract_metrics() {
|
||||
local log_file="$1"
|
||||
|
||||
local average_line reqs_line throughput reqps latency req_p90 req_p99 reqps_num
|
||||
local average_line request_line throughput reqps latency req_p90 req_p99 reqps_num
|
||||
average_line="$(extract_report_line '^[[:space:]]*[*][[:space:]]+Average:' "$log_file")"
|
||||
reqs_line="$(extract_report_line '^[[:space:]]*[*][[:space:]]+Reqs:' "$log_file")"
|
||||
request_line="$(extract_report_line '^[[:space:]]*[*][[:space:]]+(Reqs:[[:space:]]+)?Avg:' "$log_file")"
|
||||
|
||||
if [[ -n "$average_line" ]]; then
|
||||
throughput="$(echo "$average_line" | sed -E 's/^.*Average:[[:space:]]*//; s/,[[:space:]]*.*$//')"
|
||||
@@ -715,20 +715,16 @@ extract_metrics() {
|
||||
reqps="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(obj/s|req/s|ops/s|requests/s)' "$log_file")"
|
||||
fi
|
||||
|
||||
if [[ -n "$reqs_line" ]]; then
|
||||
latency="$(echo "$reqs_line" | rg -o 'Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^Avg:[[:space:]]+//')"
|
||||
req_p90="$(echo "$reqs_line" | rg -o '90%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^90%:[[:space:]]+//')"
|
||||
req_p99="$(echo "$reqs_line" | rg -o '99%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^99%:[[:space:]]+//')"
|
||||
if [[ -n "$request_line" ]]; then
|
||||
latency="$(echo "$request_line" | rg -o 'Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^Avg:[[:space:]]+//')"
|
||||
req_p90="$(echo "$request_line" | rg -o '90%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^90%:[[:space:]]+//')"
|
||||
req_p99="$(echo "$request_line" | rg -o '99%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^99%:[[:space:]]+//')"
|
||||
else
|
||||
latency="$(rg -o 'Reqs:[[:space:]]+Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^Reqs:[[:space:]]+Avg:[[:space:]]+//')"
|
||||
req_p90="$(rg -o '90%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^90%:[[:space:]]+//')"
|
||||
req_p99="$(rg -o '99%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^99%:[[:space:]]+//')"
|
||||
fi
|
||||
|
||||
if [[ -z "$latency" ]]; then
|
||||
latency="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(ms|us|µs|s)' "$log_file")"
|
||||
fi
|
||||
|
||||
throughput="$(trim "${throughput:-N/A}")"
|
||||
reqps="$(trim "${reqps:-N/A}")"
|
||||
latency="$(trim "${latency:-N/A}")"
|
||||
@@ -1128,6 +1124,8 @@ run_one_attempt() {
|
||||
"--concurrent" "$CONCURRENCY"
|
||||
"--duration" "$DURATION"
|
||||
"--region" "$REGION"
|
||||
"--no-color"
|
||||
"--analyze.v"
|
||||
)
|
||||
if [[ "$INSECURE" == "true" ]]; then
|
||||
cmd+=("--insecure")
|
||||
@@ -1212,6 +1210,12 @@ run_one_attempt() {
|
||||
req_p99_ms="$(to_ms "$req_p99_human")"
|
||||
fi
|
||||
|
||||
if [[ "$DRY_RUN" != "true" && "$TOOL" == "warp" && "$status" == "ok" ]] \
|
||||
&& rg -q '^[[:space:]]*(Total[[:space:]]+)?Errors:[[:space:]]+[1-9][0-9]*[.]?([[:space:]]|$)' "$log_file"; then
|
||||
status="failed"
|
||||
exit_code=1
|
||||
fi
|
||||
|
||||
if [[ "$DRY_RUN" != "true" && "$status" == "ok" ]]; then
|
||||
if [[ "$throughput_bps" == "N/A" && "$reqps" == "N/A" ]]; then
|
||||
status="failed"
|
||||
@@ -1318,10 +1322,21 @@ compare_baseline() {
|
||||
|
||||
dr="N/A"; dl="N/A"; dt="N/A"; dp90="N/A"; dp99="N/A"; ne="N/A"; be="N/A"; de="N/A"
|
||||
if (br!="N/A" && n_req!="N/A" && br+0!=0) dr=sprintf("%.2f", ((n_req-br)/br)*100)
|
||||
if (bl!="N/A" && n_lat!="N/A" && bl+0!=0) dl=sprintf("%.2f", ((n_lat-bl)/bl)*100)
|
||||
if (bl!="N/A" && n_lat!="N/A") {
|
||||
if (bl+0!=0) dl=sprintf("%.2f", ((n_lat-bl)/bl)*100)
|
||||
else if (n_lat+0==0) dl="0.00"
|
||||
}
|
||||
if (bt!="N/A" && n_thr!="N/A" && bt+0!=0) dt=sprintf("%.2f", ((n_thr-bt)/bt)*100)
|
||||
if (bp90!="N/A" && n_p90!="N/A" && bp90+0!=0) dp90=sprintf("%.2f", ((n_p90-bp90)/bp90)*100)
|
||||
if (bp99!="N/A" && n_p99!="N/A" && bp99+0!=0) dp99=sprintf("%.2f", ((n_p99-bp99)/bp99)*100)
|
||||
# Warp v1 rounds sub-millisecond latency to 0s. Two zero readings are
|
||||
# the same below-resolution bucket; a nonzero candidate remains invalid.
|
||||
if (bp90!="N/A" && n_p90!="N/A") {
|
||||
if (bp90+0!=0) dp90=sprintf("%.2f", ((n_p90-bp90)/bp90)*100)
|
||||
else if (n_p90+0==0) dp90="0.00"
|
||||
}
|
||||
if (bp99!="N/A" && n_p99!="N/A") {
|
||||
if (bp99+0!=0) dp99=sprintf("%.2f", ((n_p99-bp99)/bp99)*100)
|
||||
else if (n_p99+0==0) dp99="0.00"
|
||||
}
|
||||
if (n_ok!="N/A" && n_fail!="N/A" && n_ok+n_fail>0) ne=sprintf("%.2f", (n_fail/(n_ok+n_fail))*100)
|
||||
if (bok!="N/A" && bfail!="N/A" && bok+bfail>0) be=sprintf("%.2f", (bfail/(bok+bfail))*100)
|
||||
if (ne!="N/A" && be!="N/A") de=sprintf("%.2f", ne-be)
|
||||
|
||||
@@ -58,8 +58,13 @@ rg -qx 'evidence_mode=dry-run' "$OUT_DIR/manifest.env"
|
||||
rg -qx 'formal_evidence=false' "$OUT_DIR/manifest.env"
|
||||
rg -qx 'performance_conclusion=not_measured_dry_run' "$OUT_DIR/manifest.env"
|
||||
rg -qx 'bucket_isolation=per-leg' "$OUT_DIR/manifest.env"
|
||||
rg -qx 'dataset_setup=get-and-mixed-via-warp-put' "$OUT_DIR/manifest.env"
|
||||
[[ "$(rg -c -- '--extra-args --noclear' "$TRACE_FILE")" == "64" ]]
|
||||
rg -qx 'dataset_setup=get-and-mixed-via-bounded-warp-native' "$OUT_DIR/manifest.env"
|
||||
rg -qx 'dataset_objects=64' "$OUT_DIR/manifest.env"
|
||||
[[ "$(rg -c -- '--extra-args --objects\\ 64\\ --noclear' "$TRACE_FILE")" == "32" ]]
|
||||
if rg -q -- 'dataset-setup' "$TRACE_FILE"; then
|
||||
echo "unexpected redundant dataset setup command" >&2
|
||||
exit 1
|
||||
fi
|
||||
! rg -q -- 'rustfs-bench' "$TRACE_FILE"
|
||||
|
||||
if "$RUNNER" \
|
||||
|
||||
@@ -74,13 +74,28 @@ FAKE_WARP="${TMP_DIR}/fake-warp"
|
||||
cat >"$FAKE_WARP" <<'EOF'
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
[[ " $* " == *" --analyze.v "* ]]
|
||||
[[ " $* " == *" --no-color "* ]]
|
||||
if [[ "${FAKE_WARP_ZERO_LATENCY:-0}" == "1" ]]; then
|
||||
cat <<'LOG'
|
||||
Operation: GET. Concurrency: 8. Ran: 7s
|
||||
Requests considered: 1000:
|
||||
* Average: 160.00 MiB/s, 40960.00 obj/s
|
||||
* Avg: 0s, 50%: 0s, 90%: 0s, 99%: 0s, Fastest: 0s, Slowest: 1ms, StdDev: 0s
|
||||
LOG
|
||||
exit 0
|
||||
fi
|
||||
cat <<'LOG'
|
||||
- PUT Average: 161 Obj/s, 5.0MiB/s; Current 161 Obj/s, 5.0MiB/s.
|
||||
Report: GET. Concurrency: 64. Ran: 7s
|
||||
Operation: GET. Concurrency: 64. Ran: 7s
|
||||
Requests considered: 1000:
|
||||
* Average: 653.90 MiB/s, 20925.58 obj/s
|
||||
* Reqs: Avg: 3.5ms, 50%: 2.0ms, 90%: 3.6ms, 99%: 24.1ms, Fastest: 0.2ms, Slowest: 607.7ms, StdDev: 20.6ms
|
||||
* Avg: 3.5ms, 50%: 2.0ms, 90%: 3.6ms, 99%: 24.1ms, Fastest: 0.2ms, Slowest: 607.7ms, StdDev: 20.6ms
|
||||
Throughput, split into 7 x 1s:
|
||||
LOG
|
||||
if [[ "${FAKE_WARP_ERRORS:-0}" == "1" ]]; then
|
||||
echo 'Total Errors: 1.'
|
||||
fi
|
||||
EOF
|
||||
chmod +x "$FAKE_WARP"
|
||||
|
||||
@@ -103,6 +118,32 @@ chmod +x "$FAKE_WARP"
|
||||
|
||||
rg -q '^32767B,warp,1,1,128,ok,0,[^,]+,[^,]+,653.90 MiB/s,685663846.400000,20925.58,3.5 ms,3.500000,[^,]+,3.6 ms,3.600000,24.1 ms,24.100000$' "${TMP_DIR}/fake-warp-run/round_results.csv"
|
||||
|
||||
cat >"${TMP_DIR}/warp-no-details.log" <<'EOF'
|
||||
warp: Starting benchmark in 3s...
|
||||
Operation: PUT. Concurrency: 8
|
||||
* Average: 2.76 MiB/s, 707.03 obj/s
|
||||
EOF
|
||||
"$RUNNER" --extract-metrics-from-log "${TMP_DIR}/warp-no-details.log" >"${TMP_DIR}/warp-no-details.csv"
|
||||
rg -qx '2.76 MiB/s,2894069.760000,707.03,N/A,N/A,N/A,N/A,N/A,N/A' "${TMP_DIR}/warp-no-details.csv"
|
||||
|
||||
if FAKE_WARP_ERRORS=1 "$RUNNER" \
|
||||
--tool warp \
|
||||
--endpoint http://127.0.0.1:9000 \
|
||||
--access-key test-access \
|
||||
--secret-key test-secret \
|
||||
--sizes 32767B \
|
||||
--rounds 1 \
|
||||
--retry-per-round 1 \
|
||||
--retry-sleep-secs 1 \
|
||||
--cooldown-secs 0 \
|
||||
--duration 1s \
|
||||
--out-dir "${TMP_DIR}/fake-warp-errors" \
|
||||
--warp-bin "$FAKE_WARP" >/dev/null 2>&1; then
|
||||
echo "expected Warp request errors to fail the benchmark" >&2
|
||||
exit 1
|
||||
fi
|
||||
rg -q ',failed,1,' "${TMP_DIR}/fake-warp-errors/round_results.csv"
|
||||
|
||||
"$RUNNER" \
|
||||
--tool warp \
|
||||
--endpoint http://127.0.0.1:9000 \
|
||||
@@ -132,4 +173,28 @@ awk -F',' '
|
||||
END { exit found ? 0 : 1 }
|
||||
' "${TMP_DIR}/fake-warp-candidate/baseline_compare.csv"
|
||||
|
||||
for leg in baseline candidate; do
|
||||
zero_args=(
|
||||
--tool warp
|
||||
--endpoint http://127.0.0.1:9000
|
||||
--access-key test-access
|
||||
--secret-key test-secret
|
||||
--sizes 4KiB
|
||||
--rounds 1
|
||||
--retry-per-round 1
|
||||
--cooldown-secs 0
|
||||
--duration 1s
|
||||
--out-dir "${TMP_DIR}/fake-warp-zero-${leg}"
|
||||
--warp-bin "$FAKE_WARP"
|
||||
)
|
||||
if [[ "$leg" == "candidate" ]]; then
|
||||
zero_args+=(--baseline-csv "${TMP_DIR}/fake-warp-zero-baseline/median_summary.csv")
|
||||
fi
|
||||
FAKE_WARP_ZERO_LATENCY=1 "$RUNNER" "${zero_args[@]}" >/dev/null 2>&1
|
||||
done
|
||||
|
||||
"${SCRIPT_DIR}/hotpath_warp_ab_gate.sh" \
|
||||
--compare-csv "${TMP_DIR}/fake-warp-zero-candidate/baseline_compare.csv" \
|
||||
--require-tail-error >/dev/null
|
||||
|
||||
echo "object batch benchmark enhanced tests passed"
|
||||
|
||||
Reference in New Issue
Block a user