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 |
@@ -1120,48 +1120,6 @@ fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta:
|
||||
*pool_meta = previous_pool_meta;
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DecommissionCancelCommit {
|
||||
previous_start_time: Option<OffsetDateTime>,
|
||||
previous_queued: bool,
|
||||
previous_last_update: OffsetDateTime,
|
||||
canceled_pool: PoolStatus,
|
||||
}
|
||||
|
||||
fn commit_decommission_cancel(pool_meta: &mut PoolMeta, idx: usize, commit: DecommissionCancelCommit) -> Result<()> {
|
||||
let pool_count = pool_meta.pools.len();
|
||||
let Some(current) = pool_meta.pools.get(idx) else {
|
||||
return Err(invalid_decommission_pool_index_error(pool_count, idx));
|
||||
};
|
||||
// A peer reload can install the saved cancel while runtime-only fields are reconstructed.
|
||||
let cancel_already_published = commit
|
||||
.canceled_pool
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && !info.complete && !info.failed && info.start_time.is_none())
|
||||
&& PersistedPoolStatus::from(current) == PersistedPoolStatus::from(&commit.canceled_pool);
|
||||
if cancel_already_published {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let matches_generation = current.id == commit.canceled_pool.id
|
||||
&& current.cmd_line == commit.canceled_pool.cmd_line
|
||||
&& current.decommission.as_ref().is_some_and(|info| {
|
||||
info.start_time == commit.previous_start_time
|
||||
&& info.queued == commit.previous_queued
|
||||
&& is_decommission_active(info.complete, info.failed, info.canceled)
|
||||
&& (commit.previous_start_time.is_some() || current.last_update == commit.previous_last_update)
|
||||
});
|
||||
if !matches_generation {
|
||||
return Err(Error::other(format!(
|
||||
"failed to publish decommission cancel for pool {idx}: operation generation changed"
|
||||
)));
|
||||
}
|
||||
|
||||
pool_meta.pools[idx] = commit.canceled_pool;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn rollback_start_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) {
|
||||
rollback_decommission_pool_meta(pool_meta, previous_pool_meta);
|
||||
}
|
||||
@@ -1612,7 +1570,7 @@ struct PersistedPoolMeta {
|
||||
pub pools: Vec<PersistedPoolStatus>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct PersistedPoolStatus {
|
||||
#[serde(rename = "id")]
|
||||
@@ -1625,7 +1583,7 @@ struct PersistedPoolStatus {
|
||||
pub decommission: Option<PersistedPoolDecommissionInfo>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct PersistedPoolDecommissionInfo {
|
||||
#[serde(rename = "startTime", with = "time::serde::rfc3339::option")]
|
||||
@@ -3046,156 +3004,14 @@ impl ECStore {
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn decommission_cancel(self: &Arc<Self>, idx: usize) -> Result<()> {
|
||||
pub async fn decommission_cancel(&self, idx: usize) -> Result<()> {
|
||||
self.decommission_cancel_with_owner(idx, None).await
|
||||
}
|
||||
|
||||
async fn decommission_cancel_for_operation(self: &Arc<Self>, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
async fn decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> {
|
||||
self.decommission_cancel_with_owner(idx, Some(owner)).await
|
||||
}
|
||||
|
||||
async fn decommission_cancel_with_owner_and_save<Save, SaveFuture>(
|
||||
self: &Arc<Self>,
|
||||
idx: usize,
|
||||
owner: Option<&DecommissionCanceler>,
|
||||
save_pool_meta: Save,
|
||||
) -> Result<()>
|
||||
where
|
||||
Save: FnOnce(PoolMeta) -> SaveFuture + Send + 'static,
|
||||
SaveFuture: Future<Output = Result<()>> + Send + 'static,
|
||||
{
|
||||
let store = self.clone();
|
||||
let owner = owner.cloned();
|
||||
// Dropping the RPC waiter detaches this task; the transaction retains
|
||||
// the store and exact owner until persistence is resolved.
|
||||
tokio::spawn(async move { store.decommission_cancel_transaction(idx, owner, save_pool_meta).await })
|
||||
.await
|
||||
.map_err(|err| Error::other(format!("decommission cancel transaction task join error: {err}")))?
|
||||
}
|
||||
|
||||
async fn decommission_cancel_transaction<Save, SaveFuture>(
|
||||
&self,
|
||||
idx: usize,
|
||||
owner: Option<DecommissionCanceler>,
|
||||
save_pool_meta: Save,
|
||||
) -> Result<()>
|
||||
where
|
||||
Save: FnOnce(PoolMeta) -> SaveFuture,
|
||||
SaveFuture: Future<Output = Result<()>>,
|
||||
{
|
||||
let owner = owner.as_ref();
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
let operation_gate = self.ctx.decommission_operation_gate();
|
||||
let operation_guard = operation_gate.write().await;
|
||||
let save_guard = self.pool_meta_save_gate.lock().await;
|
||||
|
||||
// Lock order: start gate, operation gate, save gate,
|
||||
// decommission_cancelers, then pool_meta. The state guards stay held
|
||||
// across persistence so publication and owner termination are
|
||||
// synchronous once the save reports success.
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
let (pending, should_reload_pool_meta, already_canceled, terminal_canceler) = {
|
||||
let mut already_canceled = false;
|
||||
let (pool_present, decommission_present, terminal) = if let Some(pool) = pool_meta.pools.get(idx) {
|
||||
if let Some(info) = pool.decommission.as_ref() {
|
||||
already_canceled = info.canceled;
|
||||
(
|
||||
true,
|
||||
info.has_decommission_state(),
|
||||
should_reject_decommission_cancel_as_terminal(info.complete, info.failed),
|
||||
)
|
||||
} else {
|
||||
(true, false, false)
|
||||
}
|
||||
} else {
|
||||
(false, false, false)
|
||||
};
|
||||
|
||||
ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?;
|
||||
let previous_pool = pool_meta
|
||||
.pools
|
||||
.get(idx)
|
||||
.cloned()
|
||||
.ok_or_else(|| invalid_decommission_pool_index_error(pool_meta.pools.len(), idx))?;
|
||||
let previous_decommission = previous_pool
|
||||
.decommission
|
||||
.as_ref()
|
||||
.ok_or_else(|| decommission_metadata_not_initialized_error("cancel decommission"))?;
|
||||
let mut snapshot = pool_meta.clone();
|
||||
let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut snapshot, idx, owner, |pool_meta| {
|
||||
pool_meta.decommission_cancel(idx)
|
||||
}) else {
|
||||
return Ok(());
|
||||
};
|
||||
let pending = if changed {
|
||||
let canceled_pool = snapshot
|
||||
.pools
|
||||
.get(idx)
|
||||
.cloned()
|
||||
.ok_or_else(|| invalid_decommission_pool_index_error(pool_meta.pools.len(), idx))?;
|
||||
Some((
|
||||
snapshot,
|
||||
DecommissionCancelCommit {
|
||||
previous_start_time: previous_decommission.start_time,
|
||||
previous_queued: previous_decommission.queued,
|
||||
previous_last_update: previous_pool.last_update,
|
||||
canceled_pool,
|
||||
},
|
||||
))
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
cancelers.get(idx).and_then(Option::as_ref).cloned()
|
||||
};
|
||||
(
|
||||
pending,
|
||||
should_retry_decommission_cancel_reload(changed, already_canceled),
|
||||
already_canceled,
|
||||
terminal_canceler,
|
||||
)
|
||||
};
|
||||
let active_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active);
|
||||
if !active_worker && !already_canceled {
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "cancel_skipped",
|
||||
reason = "no_active_canceler",
|
||||
"Decommission cancel skipped"
|
||||
);
|
||||
}
|
||||
|
||||
let commit_result = if let Some((snapshot, commit)) = pending {
|
||||
save_pool_meta(snapshot).await?;
|
||||
commit_decommission_cancel(&mut pool_meta, idx, commit)
|
||||
} else {
|
||||
Ok(())
|
||||
};
|
||||
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, canceler);
|
||||
}
|
||||
commit_result?;
|
||||
drop(pool_meta);
|
||||
drop(cancelers);
|
||||
drop(save_guard);
|
||||
drop(operation_guard);
|
||||
|
||||
if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("decommission_cancel for pool {idx}");
|
||||
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn release_decommission_canceler_slot(&self, idx: usize, owner: &DecommissionCanceler) {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
take_and_cancel_decommission_canceler_for_operation(cancelers.as_mut_slice(), idx, owner);
|
||||
@@ -3223,11 +3039,7 @@ impl ECStore {
|
||||
retryable
|
||||
}
|
||||
|
||||
async fn retry_decommission_cancel_for_operation(self: &Arc<Self>, idx: usize, owner: &DecommissionCanceler) {
|
||||
if !self.decommission_terminal_retryable_for_operation(idx, owner).await {
|
||||
return;
|
||||
}
|
||||
|
||||
async fn retry_decommission_cancel_for_operation(&self, idx: usize, owner: &DecommissionCanceler) {
|
||||
let mut attempt = 0usize;
|
||||
loop {
|
||||
let Err(err) = self.decommission_cancel_for_operation(idx, owner).await else {
|
||||
@@ -3277,10 +3089,87 @@ impl ECStore {
|
||||
}
|
||||
}
|
||||
|
||||
async fn decommission_cancel_with_owner(self: &Arc<Self>, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
|
||||
let pools = self.pools.clone();
|
||||
self.decommission_cancel_with_owner_and_save(idx, owner, move |snapshot| async move { snapshot.save(pools).await })
|
||||
.await
|
||||
async fn decommission_cancel_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> {
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "cancel decommission")?;
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
|
||||
// Lock order: decommission_cancelers before pool_meta. Holding both makes
|
||||
// owner validation and the terminal transition one atomic operation.
|
||||
let (should_save_pool_meta, should_reload_pool_meta, already_canceled, previous_pool_meta, terminal_canceler) = {
|
||||
let cancelers = self.decommission_cancelers.read().await;
|
||||
let mut lock = self.pool_meta.write().await;
|
||||
let mut already_canceled = false;
|
||||
let (pool_present, decommission_present, terminal) = if let Some(pool) = lock.pools.get(idx) {
|
||||
if let Some(info) = pool.decommission.as_ref() {
|
||||
already_canceled = info.canceled;
|
||||
(
|
||||
true,
|
||||
info.has_decommission_state(),
|
||||
should_reject_decommission_cancel_as_terminal(info.complete, info.failed),
|
||||
)
|
||||
} else {
|
||||
(true, false, false)
|
||||
}
|
||||
} else {
|
||||
(false, false, false)
|
||||
};
|
||||
|
||||
ensure_decommission_cancel_allowed(pool_present, decommission_present, terminal)?;
|
||||
let previous_pool_meta = lock.clone();
|
||||
let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut lock, idx, owner, |pool_meta| {
|
||||
pool_meta.decommission_cancel(idx)
|
||||
}) else {
|
||||
return Ok(());
|
||||
};
|
||||
let terminal_canceler = if let Some(owner) = owner {
|
||||
Some(owner.clone())
|
||||
} else {
|
||||
cancelers.get(idx).and_then(Option::as_ref).cloned()
|
||||
};
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
canceler.cancel();
|
||||
}
|
||||
(
|
||||
changed,
|
||||
should_retry_decommission_cancel_reload(changed, already_canceled),
|
||||
already_canceled,
|
||||
changed.then_some(previous_pool_meta),
|
||||
terminal_canceler,
|
||||
)
|
||||
};
|
||||
let canceled_worker = terminal_canceler.as_ref().is_some_and(DecommissionCanceler::is_active);
|
||||
if !canceled_worker && !already_canceled {
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "cancel_skipped",
|
||||
reason = "no_active_canceler",
|
||||
"Decommission cancel skipped"
|
||||
);
|
||||
}
|
||||
|
||||
self.wait_for_decommission_side_effects().await;
|
||||
|
||||
if should_save_pool_meta && let Err(err) = self.save_current_pool_meta().await {
|
||||
if let Some(previous_pool_meta) = previous_pool_meta {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta);
|
||||
}
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
if let Some(canceler) = terminal_canceler.as_ref() {
|
||||
self.release_decommission_canceler_slot(idx, canceler).await;
|
||||
}
|
||||
|
||||
if should_reload_pool_meta && let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("decommission_cancel for pool {idx}");
|
||||
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
@@ -9756,525 +9645,6 @@ mod pools_tests {
|
||||
assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_save_failure_preserves_generation_and_token_until_retry() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
version: super::POOL_META_VERSION,
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
|
||||
let err = store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), |_| async { Err(Error::Timeout) })
|
||||
.await
|
||||
.expect_err("injected pool metadata timeout should fail cancel");
|
||||
assert!(matches!(err, Error::Timeout));
|
||||
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("failed cancel must retain decommission metadata");
|
||||
assert_eq!(info.start_time, Some(generation));
|
||||
assert!(!info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(ensure_decommission_generation(&pool_meta, 0, generation).is_ok());
|
||||
}
|
||||
{
|
||||
let cancelers = store.decommission_cancelers.read().await;
|
||||
let current = cancelers[0].as_ref().expect("failed cancel must retain the worker owner");
|
||||
assert!(current.owns_same_operation(&canceler));
|
||||
assert!(current.is_active());
|
||||
}
|
||||
assert!(!canceler.is_cancelled());
|
||||
assert!(store.decommission_terminal_retryable_for_operation(0, &canceler).await);
|
||||
|
||||
let persisted = Arc::new(std::sync::Mutex::new(None));
|
||||
let persisted_for_save = persisted.clone();
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
let data = snapshot.encode_config_data()?;
|
||||
*persisted_for_save
|
||||
.lock()
|
||||
.expect("persisted snapshot lock should not be poisoned") = Some(data);
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.expect("cancel retry should commit");
|
||||
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("committed cancel metadata should remain present");
|
||||
assert!(info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(info.start_time.is_none());
|
||||
}
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
|
||||
let repeated_save_called = Arc::new(AtomicBool::new(false));
|
||||
let repeated_save_called_by_closure = repeated_save_called.clone();
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, None, move |_| async move {
|
||||
repeated_save_called_by_closure.store(true, Ordering::SeqCst);
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.expect("repeated cancel should be idempotent");
|
||||
assert!(!repeated_save_called.load(Ordering::SeqCst));
|
||||
|
||||
let persisted = persisted
|
||||
.lock()
|
||||
.expect("persisted snapshot lock should not be poisoned")
|
||||
.take()
|
||||
.expect("successful retry should capture persisted bytes");
|
||||
let mut restarted = PoolMeta::default();
|
||||
restarted
|
||||
.load_from_config_data(persisted)
|
||||
.expect("a restarted process should decode the committed cancel");
|
||||
assert!(
|
||||
restarted.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && info.start_time.is_none())
|
||||
);
|
||||
assert!(first_resumable_decommission_queue_indices(&restarted).is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_serializes_reload_until_local_commit() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
version: super::POOL_META_VERSION,
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
stage: "migrate_object".to_string(),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let (persisted_tx, persisted_rx) = tokio::sync::oneshot::channel();
|
||||
let save_release = Arc::new(tokio::sync::Notify::new());
|
||||
|
||||
let cancel = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
let save_release = save_release.clone();
|
||||
async move {
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
persisted_tx
|
||||
.send(snapshot.encode_config_data()?)
|
||||
.map_err(|_| Error::other("failed to expose saved cancel snapshot"))?;
|
||||
save_release.notified().await;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
});
|
||||
let persisted = persisted_rx.await.expect("save should expose the canceled snapshot");
|
||||
let reload_started = Arc::new(tokio::sync::Notify::new());
|
||||
let reload = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let reload_started = reload_started.clone();
|
||||
async move {
|
||||
let mut reloaded = PoolMeta::default();
|
||||
reloaded.load_from_config_data(persisted)?;
|
||||
reload_started.notify_one();
|
||||
*store.pool_meta.write().await = reloaded;
|
||||
Ok::<(), Error>(())
|
||||
}
|
||||
});
|
||||
reload_started.notified().await;
|
||||
assert!(!reload.is_finished(), "peer reload must wait for cancel publication");
|
||||
|
||||
save_release.notify_one();
|
||||
cancel
|
||||
.await
|
||||
.expect("cancel task should not panic")
|
||||
.expect("cancel should commit before releasing the reload");
|
||||
reload
|
||||
.await
|
||||
.expect("reload task should not panic")
|
||||
.expect("reload should install the saved cancel after publication");
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("reloaded cancel metadata should remain present");
|
||||
assert!(info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(info.start_time.is_none());
|
||||
assert!(info.stage.is_empty(), "the persisted snapshot should have been decoded before commit");
|
||||
drop(pool_meta);
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_transaction_survives_caller_abort_after_durable_save() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let persisted = Arc::new(std::sync::Mutex::new(None));
|
||||
let (durable_tx, durable_rx) = tokio::sync::oneshot::channel();
|
||||
let save_release = Arc::new(tokio::sync::Notify::new());
|
||||
let cancel = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
let persisted = persisted.clone();
|
||||
let save_release = save_release.clone();
|
||||
async move {
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
assert!(
|
||||
snapshot.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && info.start_time.is_none())
|
||||
);
|
||||
*persisted.lock().expect("persisted cancel lock should not be poisoned") =
|
||||
Some(snapshot.encode_config_data()?);
|
||||
durable_tx
|
||||
.send(())
|
||||
.map_err(|_| Error::other("failed to report durable cancel snapshot"))?;
|
||||
save_release.notified().await;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
});
|
||||
durable_rx.await.expect("save hook should report the durable cancel snapshot");
|
||||
cancel.abort();
|
||||
let join_err = cancel.await.expect_err("caller cancel future should be aborted");
|
||||
assert!(join_err.is_cancelled());
|
||||
assert!(
|
||||
store.pool_meta.try_read().is_err(),
|
||||
"the detached transaction must retain its state guard"
|
||||
);
|
||||
assert!(
|
||||
store.decommission_cancelers.try_read().is_err(),
|
||||
"the detached transaction must retain its owner guard"
|
||||
);
|
||||
|
||||
save_release.notify_one();
|
||||
tokio::time::timeout(StdDuration::from_secs(1), canceler.token().cancelled())
|
||||
.await
|
||||
.expect("detached cancel transaction should terminate the old token");
|
||||
|
||||
let persisted = persisted
|
||||
.lock()
|
||||
.expect("persisted cancel lock should not be poisoned")
|
||||
.take()
|
||||
.expect("save should capture the durable cancel snapshot");
|
||||
let mut durable = PoolMeta::default();
|
||||
durable
|
||||
.load_from_config_data(persisted)
|
||||
.expect("durable cancel snapshot should decode");
|
||||
assert!(
|
||||
durable.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && info.start_time.is_none())
|
||||
);
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
|
||||
let side_effect_ran = Arc::new(AtomicBool::new(false));
|
||||
let operation_gate = store.ctx.decommission_operation_gate();
|
||||
let result = run_decommission_side_effect(canceler.token(), &operation_gate, {
|
||||
let side_effect_ran = side_effect_ran.clone();
|
||||
move || async move {
|
||||
side_effect_ran.store(true, Ordering::SeqCst);
|
||||
Ok(())
|
||||
}
|
||||
})
|
||||
.await;
|
||||
assert!(matches!(result, Err(Error::OperationCanceled)));
|
||||
assert!(!side_effect_ran.load(Ordering::SeqCst));
|
||||
assert!(
|
||||
store.pool_meta.read().await.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.is_some_and(|info| info.canceled && !info.failed)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_quiesces_and_publishes_only_after_save() {
|
||||
let generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let operation_gate = store.ctx.decommission_operation_gate();
|
||||
let side_effect = operation_gate.read().await;
|
||||
let save_started = Arc::new(tokio::sync::Notify::new());
|
||||
let save_release = Arc::new(tokio::sync::Notify::new());
|
||||
let save_entered = Arc::new(AtomicBool::new(false));
|
||||
|
||||
let mut cancel = tokio::spawn({
|
||||
let store = store.clone();
|
||||
let canceler = canceler.clone();
|
||||
let save_started = save_started.clone();
|
||||
let save_release = save_release.clone();
|
||||
let save_entered = save_entered.clone();
|
||||
async move {
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |_| async move {
|
||||
save_entered.store(true, Ordering::SeqCst);
|
||||
save_started.notify_one();
|
||||
save_release.notified().await;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
});
|
||||
|
||||
let mut cancel_holds_start_gate = false;
|
||||
for _ in 0..100 {
|
||||
if store.start_gate.try_lock().is_err() {
|
||||
cancel_holds_start_gate = true;
|
||||
break;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
assert!(cancel_holds_start_gate, "cancel should reach the operation-gate wait");
|
||||
assert!(!cancel.is_finished(), "cancel must wait for the in-flight side effect");
|
||||
assert!(!save_entered.load(Ordering::SeqCst), "cancel must quiesce side effects before saving");
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
assert!(ensure_decommission_generation(&pool_meta, 0, generation).is_ok());
|
||||
}
|
||||
assert!(!canceler.is_cancelled());
|
||||
|
||||
drop(side_effect);
|
||||
tokio::time::timeout(StdDuration::from_secs(1), save_started.notified())
|
||||
.await
|
||||
.expect("cancel should reach the injected save after side effects quiesce");
|
||||
assert!(
|
||||
store.pool_meta_save_gate.try_lock().is_err(),
|
||||
"the cancel save must exclude stale full-document saves until publication"
|
||||
);
|
||||
assert!(
|
||||
store.pool_meta.try_read().is_err(),
|
||||
"the active generation must stay write-locked through persistence"
|
||||
);
|
||||
assert!(!canceler.is_cancelled(), "the token must remain live until persistence commits");
|
||||
|
||||
let mut complete = tokio::spawn({
|
||||
let store = store.clone();
|
||||
async move { store.complete_decommission(0).await }
|
||||
});
|
||||
let mut fail = tokio::spawn({
|
||||
let store = store.clone();
|
||||
async move { store.decommission_failed(0).await }
|
||||
});
|
||||
tokio::task::yield_now().await;
|
||||
assert!(!complete.is_finished(), "complete must serialize behind the pending cancel");
|
||||
assert!(!fail.is_finished(), "fail must serialize behind the pending cancel");
|
||||
|
||||
save_release.notify_one();
|
||||
tokio::time::timeout(StdDuration::from_secs(1), &mut cancel)
|
||||
.await
|
||||
.expect("cancel should finish after persistence commits")
|
||||
.expect("cancel task should not panic")
|
||||
.expect("cancel should commit");
|
||||
tokio::time::timeout(StdDuration::from_secs(1), &mut complete)
|
||||
.await
|
||||
.expect("complete should finish after cancel commits")
|
||||
.expect("complete task should not panic")
|
||||
.expect("stale complete should be a no-op");
|
||||
tokio::time::timeout(StdDuration::from_secs(1), &mut fail)
|
||||
.await
|
||||
.expect("fail should finish after cancel commits")
|
||||
.expect("fail task should not panic")
|
||||
.expect("stale fail should be a no-op");
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("cancel metadata should remain present");
|
||||
assert!(info.canceled);
|
||||
assert!(!info.complete);
|
||||
assert!(!info.failed);
|
||||
assert!(info.start_time.is_none());
|
||||
drop(pool_meta);
|
||||
assert!(canceler.is_cancelled());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_decommission_cancel_commit_rejects_a_replaced_generation() {
|
||||
let old_generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let new_generation = old_generation + Duration::seconds(1);
|
||||
let mut canceled = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(old_generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let previous_last_update = canceled.pools[0].last_update;
|
||||
assert!(canceled.decommission_cancel(0));
|
||||
let canceled_pool = canceled.pools.remove(0);
|
||||
let commit = super::DecommissionCancelCommit {
|
||||
previous_start_time: Some(old_generation),
|
||||
previous_queued: false,
|
||||
previous_last_update,
|
||||
canceled_pool,
|
||||
};
|
||||
let mut current = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(new_generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let err = super::commit_decommission_cancel(&mut current, 0, commit)
|
||||
.expect_err("an old cancel must not publish over a replacement generation");
|
||||
assert!(err.to_string().contains("operation generation changed"));
|
||||
let info = current.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("replacement generation should remain present");
|
||||
assert_eq!(info.start_time, Some(new_generation));
|
||||
assert!(!info.canceled);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_cancel_rejects_stale_retry_after_queued_replacement() {
|
||||
let old_generation = OffsetDateTime::UNIX_EPOCH;
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
let pool_meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(old_generation),
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
let store = decommission_worker_test_store(pool_meta, vec![Some(canceler.clone())]);
|
||||
let queued_replacement = Arc::new(std::sync::Mutex::new(None));
|
||||
let queued_replacement_for_save = queued_replacement.clone();
|
||||
|
||||
store
|
||||
.decommission_cancel_with_owner_and_save(0, Some(&canceler), move |snapshot| async move {
|
||||
let saved_cancel = snapshot.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("saved snapshot should contain decommission metadata");
|
||||
assert!(saved_cancel.canceled);
|
||||
assert!(saved_cancel.start_time.is_none());
|
||||
let mut queued_replacement = snapshot.clone();
|
||||
let replacement = queued_replacement
|
||||
.pools
|
||||
.get_mut(0)
|
||||
.expect("cancel snapshot should contain the pool");
|
||||
replacement.last_update += Duration::seconds(1);
|
||||
let info = replacement
|
||||
.decommission
|
||||
.as_mut()
|
||||
.expect("cancel snapshot should contain decommission metadata");
|
||||
info.canceled = false;
|
||||
info.queued = true;
|
||||
*queued_replacement_for_save
|
||||
.lock()
|
||||
.expect("queued replacement lock should not be poisoned") = Some(queued_replacement);
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.expect("cancel should commit before a queued replacement is installed");
|
||||
assert!(store.decommission_cancelers.read().await[0].is_none());
|
||||
assert!(!canceler.is_active());
|
||||
assert!(canceler.is_cancelled());
|
||||
|
||||
let queued_replacement = queued_replacement
|
||||
.lock()
|
||||
.expect("queued replacement lock should not be poisoned")
|
||||
.take()
|
||||
.expect("save should prepare the queued replacement");
|
||||
*store.pool_meta.write().await = queued_replacement;
|
||||
|
||||
let replacement_revision = {
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("queued replacement should remain present");
|
||||
assert!(info.queued);
|
||||
assert!(!info.canceled);
|
||||
assert!(info.start_time.is_none());
|
||||
pool_meta.pools[0].last_update
|
||||
};
|
||||
|
||||
store.retry_decommission_cancel_for_operation(0, &canceler).await;
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("stale retry must preserve the queued replacement");
|
||||
assert_eq!(pool_meta.pools[0].last_update, replacement_revision);
|
||||
assert!(info.queued);
|
||||
assert!(!info.canceled);
|
||||
assert!(info.start_time.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_decommission_failed_save_failure_preserves_owner_until_retry_succeeds() {
|
||||
let canceler = DecommissionCanceler::new(CancellationToken::new());
|
||||
|
||||
@@ -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