Compare commits

..

3 Commits

Author SHA1 Message Date
overtrue 82bff87f19 fix(ci): accept equal below-resolution Warp latency 2026-08-23 09:54:40 +08:00
overtrue 78b6879043 fix(ci): collect valid Warp tail and error evidence 2026-08-23 07:39:25 +08:00
overtrue 520b93c5fc fix(ci): bound Warp ABBA dataset preparation 2026-08-23 05:00:47 +08:00
5 changed files with 200 additions and 760 deletions
+86 -716
View File
@@ -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());
+11 -26
View File
@@ -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=""
+29 -14
View File
@@ -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)
+7 -2
View File
@@ -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" \
+67 -2
View File
@@ -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"