perf(ilm): reduce transition transaction mutations (#7320)

* perf(ilm): reduce transition transaction mutations

* test(ilm): rename transition kill points

---------

Co-authored-by: Zhengchao An <anzhengchao@gmail.com>
This commit is contained in:
cxymds
2026-09-07 10:40:57 +08:00
committed by GitHub
parent d633a635ec
commit 6018dd372f
11 changed files with 1135 additions and 128 deletions
@@ -499,9 +499,7 @@ impl DurableIlmRecordCheckpoint {
},
) => {
previous_identity == next_identity
&& transition_state_distance(*previous_state, *next_state)
.and_then(|distance| previous_revision.checked_add(distance))
.is_some_and(|expected_revision| *next_revision == expected_revision)
&& transition_state_revision_is_successor(*previous_state, *previous_revision, *next_state, *next_revision)
&& (!previous_remote_version_known || previous_remote_version == next_remote_version)
}
(
@@ -1030,6 +1028,22 @@ fn transition_state_distance(
}
}
fn transition_state_revision_is_successor(
from: transition_transaction::TransitionTransactionState,
from_revision: u64,
to: transition_transaction::TransitionTransactionState,
to_revision: u64,
) -> bool {
use transition_transaction::TransitionTransactionState::{LocalCommitStarted, UploadOutcomeUnknown};
if from == UploadOutcomeUnknown && from_revision == 1 && to == LocalCommitStarted {
return to_revision == 2;
}
transition_state_distance(from, to)
.and_then(|distance| from_revision.checked_add(distance))
.is_some_and(|expected_revision| to_revision == expected_revision)
}
fn tier_probe_state_reaches(
from: tier_probe_intent::TierProbeIntentState,
to: tier_probe_intent::TierProbeIntentState,
@@ -2173,6 +2187,64 @@ mod tests {
);
}
#[test]
fn transition_checkpoint_accepts_only_the_distinguishable_compact_edge() {
let identity_sha256 = "a".repeat(64);
let unknown_remote_sha256 = "b".repeat(64);
let known_remote_sha256 = "c".repeat(64);
let checkpoint = |revision, state, remote_version_sha256: String, remote_version_known| {
DurableIlmRecordCheckpoint::TransitionTransaction {
content_sha256: format!("{revision:064x}"),
identity_sha256: identity_sha256.clone(),
remote_version_sha256,
remote_version_known,
revision,
state,
}
};
let compact_unknown = checkpoint(
1,
transition_transaction::TransitionTransactionState::UploadOutcomeUnknown,
unknown_remote_sha256.clone(),
false,
);
let compact_local_commit = checkpoint(
2,
transition_transaction::TransitionTransactionState::LocalCommitStarted,
known_remote_sha256.clone(),
true,
);
compact_unknown
.validate_successor(&compact_local_commit)
.expect("compact pre-upload fence should advance directly to the exact local-commit fence");
let legacy_unknown = checkpoint(
2,
transition_transaction::TransitionTransactionState::UploadOutcomeUnknown,
unknown_remote_sha256,
false,
);
let invalid_legacy_skip = checkpoint(
3,
transition_transaction::TransitionTransactionState::LocalCommitStarted,
known_remote_sha256.clone(),
true,
);
assert!(
legacy_unknown.validate_successor(&invalid_legacy_skip).is_err(),
"legacy UploadOutcomeUnknown@2 must not masquerade as the compact edge"
);
let valid_legacy_skip = checkpoint(
4,
transition_transaction::TransitionTransactionState::LocalCommitStarted,
known_remote_sha256,
true,
);
legacy_unknown
.validate_successor(&valid_legacy_skip)
.expect("legacy receipts may still observe the existing two-edge state advance");
}
#[test]
fn tier_delete_dispatch_manifest_namespace_validates_monotonic_branches() {
use tier_delete_journal::TierDeleteDispatchManifestState::{Aborted, Aborting, Completed, DispatchAuthorized, Preparing};
@@ -273,6 +273,14 @@ pub struct TransitionTransactionInit {
impl TransitionTransaction {
pub fn new(init: TransitionTransactionInit) -> Result<Self> {
Self::new_with_initial_state(init, TransitionTransactionState::UploadStarted)
}
pub(crate) fn new_compact(init: TransitionTransactionInit) -> Result<Self> {
Self::new_with_initial_state(init, TransitionTransactionState::UploadOutcomeUnknown)
}
fn new_with_initial_state(init: TransitionTransactionInit, state: TransitionTransactionState) -> Result<Self> {
let remote_object =
canonical_transition_remote_object(init.deployment_id, &init.source.bucket, init.transaction_id, init.write_id)?;
let transaction = Self {
@@ -286,7 +294,7 @@ impl TransitionTransaction {
backend_fingerprint: init.backend_fingerprint,
remote_object,
remote_version: TransitionRemoteVersion::unknown(),
state: TransitionTransactionState::UploadStarted,
state,
not_after_unix_nanos: init.not_after_unix_nanos,
};
transaction.validate()?;
@@ -356,7 +364,7 @@ impl TransitionTransaction {
remote_version: Option<TransitionRemoteVersion>,
) -> Result<TransitionTransactionFence> {
self.check_fence(fence)?;
if !state_change_allowed(self.state, next) {
if !state_change_allowed_at(self.state, next, self.revision) {
return Err(TransitionTransactionError::InvalidStateChange {
from: self.state,
to: next,
@@ -388,6 +396,14 @@ impl TransitionTransaction {
}
self.remote_version = TransitionRemoteVersion::unknown();
}
TransitionTransactionState::LocalCommitStarted if self.state == TransitionTransactionState::UploadOutcomeUnknown => {
let remote_version =
remote_version.ok_or(TransitionTransactionError::Corrupt("compact local commit requires remote version"))?;
if remote_version.is_unknown() {
return Err(TransitionTransactionError::Corrupt("compact local commit requires known remote version"));
}
self.remote_version = remote_version;
}
TransitionTransactionState::LocalCommitStarted | TransitionTransactionState::Committed => {
if let Some(remote_version) = remote_version
&& remote_version != self.remote_version
@@ -640,7 +656,7 @@ pub(crate) async fn save_transition_transaction_record_if_current(
) -> EcstoreResult<()> {
let object = transition_transaction_record_object_name(next.transaction_id).map_err(transition_transaction_store_error)?;
let revision_is_next = expected.revision.checked_add(1) == Some(next.revision);
let state_is_next = state_change_allowed(expected.state, next.state)
let state_is_next = state_change_allowed_at(expected.state, next.state, expected.revision)
|| matches!(
(expected.state, next.state),
(
@@ -2146,7 +2162,7 @@ where
}
}
fn state_change_allowed(from: TransitionTransactionState, to: TransitionTransactionState) -> bool {
fn state_change_allowed_at(from: TransitionTransactionState, to: TransitionTransactionState, revision: u64) -> bool {
matches!(
(from, to),
(TransitionTransactionState::UploadStarted, TransitionTransactionState::Uploaded)
@@ -2158,7 +2174,9 @@ fn state_change_allowed(from: TransitionTransactionState, to: TransitionTransact
| (TransitionTransactionState::UploadOutcomeUnknown, TransitionTransactionState::Uploaded)
| (TransitionTransactionState::Uploaded, TransitionTransactionState::LocalCommitStarted)
| (TransitionTransactionState::LocalCommitStarted, TransitionTransactionState::Committed)
)
) || (revision == 1
&& from == TransitionTransactionState::UploadOutcomeUnknown
&& to == TransitionTransactionState::LocalCommitStarted)
}
fn state_requires_known_remote_version(state: TransitionTransactionState) -> bool {
@@ -2611,6 +2629,57 @@ mod tests {
assert_eq!(transaction.state, TransitionTransactionState::Committed);
}
#[test]
fn compact_state_sequence_is_distinguishable_and_keeps_legacy_edges_strict() {
let init = TransitionTransactionInit {
deployment_id: Uuid::new_v4(),
transaction_id: Uuid::new_v4(),
owner_epoch: Uuid::new_v4(),
write_id: Uuid::new_v4(),
source: source_identity(TransitionSourceVersionMode::Versioned),
tier_name: "warm-tier".to_string(),
backend_fingerprint: BACKEND_FINGERPRINT,
not_after_unix_nanos: 1_780_000_000_000_000_000,
};
let mut compact = TransitionTransaction::new_compact(init).expect("compact transaction should be created");
assert_eq!(compact.state, TransitionTransactionState::UploadOutcomeUnknown);
assert_eq!(compact.revision, 1);
let remote_version = TransitionRemoteVersion::versioned(Uuid::new_v4().to_string());
let fence = compact
.advance(
compact.fence(),
TransitionTransactionState::LocalCommitStarted,
Some(remote_version.clone()),
)
.expect("compact upload should persist its exact candidate at the local commit fence");
assert_eq!(fence.revision, 2);
assert_eq!(compact.remote_version, remote_version);
assert_eq!(compact.state, TransitionTransactionState::LocalCommitStarted);
let encoded = compact
.encode()
.expect("compact transaction should encode as v1-compatible bytes");
assert_eq!(
TransitionTransaction::decode(compact.transaction_id, &encoded).expect("compact transaction should decode"),
compact
);
let mut legacy_unknown = new_transaction();
legacy_unknown
.advance(legacy_unknown.fence(), TransitionTransactionState::UploadOutcomeUnknown, None)
.expect("legacy transaction should persist its pre-upload fence");
assert!(matches!(
legacy_unknown.advance(
legacy_unknown.fence(),
TransitionTransactionState::LocalCommitStarted,
Some(TransitionRemoteVersion::unversioned()),
),
Err(TransitionTransactionError::InvalidStateChange {
from: TransitionTransactionState::UploadOutcomeUnknown,
to: TransitionTransactionState::LocalCommitStarted,
})
));
}
#[test]
fn cleanup_pending_requires_exact_proof_and_state_specific_decision() {
let mut transaction = new_transaction();
@@ -1808,6 +1808,15 @@ impl PeerRestClient {
Ok((self.topology_member.clone(), epoch))
}
pub async fn probe_transition_transaction_compaction(&self, topology_fingerprint: String) -> Result<(String, Uuid)> {
let probe = rustfs_protos::transition_transaction_compaction_capability_probe(Uuid::new_v4().as_bytes());
let result = self
.heal_control(rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION, topology_fingerprint, probe)
.await?;
let epoch = decode_remote_version_state_capability(&self.topology_member, &result)?;
Ok((self.topology_member.clone(), epoch))
}
pub async fn load_bucket_metadata(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> {
let result = tokio::time::timeout(BUCKET_METADATA_RELOAD_TIMEOUT, async {
let result = self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await;
+130 -2
View File
@@ -318,12 +318,21 @@ pub struct IlmRecoveryExportFleetProofToken {
_permit: FleetCapabilityProofPermit,
}
/// Effect-window authority for emitting the compact transition-transaction
/// state sequence. The generation permit prevents a successor proof from
/// being published until the admitted writer has finished.
pub(crate) struct TransitionTransactionCompactionFleetProofToken {
token: FleetCapabilityProofToken,
_permit: FleetCapabilityProofPermit,
}
static REMOTE_VERSION_STATE_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static CROSS_POOL_FENCE_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static TIER_DELETE_JOURNAL_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static DECOMMISSION_TARGET_FENCE_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static LEGACY_TRANSITION_STATE_RECONCILE_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static ILM_RECOVERY_EXPORT_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static TRANSITION_TRANSACTION_COMPACTION_FLEET_PROOF: OnceLock<std::sync::RwLock<FleetCapabilityProofState>> = OnceLock::new();
static REMOTE_VERSION_STATE_PROBE_TOPOLOGY: OnceLock<String> = OnceLock::new();
static ILM_RECOVERY_EXPORT_LOCAL_PROCESS_EPOCH: LazyLock<Uuid> = LazyLock::new(Uuid::new_v4);
@@ -351,6 +360,10 @@ fn ilm_recovery_export_fleet_proof_slot() -> &'static std::sync::RwLock<FleetCap
ILM_RECOVERY_EXPORT_FLEET_PROOF.get_or_init(|| std::sync::RwLock::new(FleetCapabilityProofState::default()))
}
fn transition_transaction_compaction_fleet_proof_slot() -> &'static std::sync::RwLock<FleetCapabilityProofState> {
TRANSITION_TRANSACTION_COMPACTION_FLEET_PROOF.get_or_init(|| std::sync::RwLock::new(FleetCapabilityProofState::default()))
}
fn revoke_fleet_capability_proof_state(state: &mut FleetCapabilityProofState) {
if let Some(proof) = state.proof.take() {
proof.generation.revoke();
@@ -451,6 +464,33 @@ pub(crate) fn remote_version_state_fleet_proof_matches(proof: &RemoteVersionStat
fleet_capability_proof_matches(remote_version_state_fleet_proof_slot(), &proof.0)
}
pub(crate) fn acquire_transition_transaction_compaction_fleet_proof() -> Option<TransitionTransactionCompactionFleetProofToken> {
let expected_topology = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.get()?;
let state = transition_transaction_compaction_fleet_proof_slot()
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let token = acquire_fleet_capability_proof_from(&state, expected_topology, Instant::now())?;
let permit = state.proof.as_ref()?.generation.try_acquire()?;
Some(TransitionTransactionCompactionFleetProofToken { token, _permit: permit })
}
pub(crate) fn transition_transaction_compaction_fleet_proof_matches(
proof: &TransitionTransactionCompactionFleetProofToken,
) -> bool {
let Some(expected_topology) = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.get() else {
return false;
};
let state = transition_transaction_compaction_fleet_proof_slot()
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
proof._permit.generation.is_accepting()
&& fleet_capability_proof_matches_at(&state, &proof.token, expected_topology, Instant::now())
&& state
.proof
.as_ref()
.is_some_and(|current| Arc::ptr_eq(&current.generation, &proof._permit.generation))
}
pub fn acquire_cross_pool_fence_fleet_proof() -> Option<CrossPoolFenceFleetProofToken> {
let expected_topology = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.get()?;
let state = cross_pool_fence_fleet_proof_slot()
@@ -1073,6 +1113,35 @@ pub(crate) fn install_remote_version_state_fleet_proof_for_test(topology_fingerp
RemoteVersionStateFleetProofGuard
}
#[cfg(all(test, feature = "test-util"))]
pub(crate) struct TransitionTransactionCompactionFleetProofGuard;
#[cfg(all(test, feature = "test-util"))]
impl Drop for TransitionTransactionCompactionFleetProofGuard {
fn drop(&mut self) {
revoke_fleet_capability_proof(transition_transaction_compaction_fleet_proof_slot());
}
}
#[cfg(all(test, feature = "test-util"))]
pub(crate) fn install_transition_transaction_compaction_fleet_proof_for_test(
topology_fingerprint: &str,
) -> TransitionTransactionCompactionFleetProofGuard {
let _ = REMOTE_VERSION_STATE_PROBE_TOPOLOGY.set(topology_fingerprint.to_string());
let effective_topology = REMOTE_VERSION_STATE_PROBE_TOPOLOGY
.get()
.expect("transition transaction compaction test topology should be initialized");
if let Some(err) = publish_fleet_capability_probe_result(
transition_transaction_compaction_fleet_proof_slot(),
effective_topology,
Ok(BTreeMap::new()),
Instant::now(),
) {
panic!("test proof installation must not fail: {err}");
}
TransitionTransactionCompactionFleetProofGuard
}
fn insert_remote_version_state_peer(peer_epochs: &mut BTreeMap<String, Uuid>, peer: String, epoch: Uuid) -> Result<()> {
if epoch.is_nil() || peer_epochs.values().any(|existing| *existing == epoch) || peer_epochs.insert(peer, epoch).is_some() {
return Err(Error::other("remote version state capability peer identity is invalid"));
@@ -1090,6 +1159,7 @@ pub fn start_remote_version_state_fleet_probe(topology_fingerprint: String) {
decommission_target_fence_fleet_proof_slot(),
legacy_transition_state_reconcile_fleet_proof_slot(),
ilm_recovery_export_fleet_proof_slot(),
transition_transaction_compaction_fleet_proof_slot(),
] {
mark_fleet_capability_topology_conflict(slot);
}
@@ -1133,8 +1203,25 @@ pub fn start_remote_version_state_fleet_probe(topology_fingerprint: String) {
None => Err(Error::other("ILM recovery export fleet capability notification system is unavailable")),
}
};
let (result, fence_probe, recovery_export_result) =
tokio::join!(remote_version_state_probe, cross_pool_fence_probe, recovery_export_probe);
let transition_transaction_compaction_probe = async {
match notification_sys.as_ref() {
Some(notification_sys) => timeout(
REMOTE_VERSION_STATE_PROBE_TIMEOUT,
notification_sys.probe_transition_transaction_compaction_fleet(&topology_fingerprint),
)
.await
.unwrap_or_else(|_| Err(Error::other("transition transaction compaction fleet capability probe timed out"))),
None => Err(Error::other(
"transition transaction compaction fleet capability notification system is unavailable",
)),
}
};
let (result, fence_probe, recovery_export_result, transition_transaction_compaction_result) = tokio::join!(
remote_version_state_probe,
cross_pool_fence_probe,
recovery_export_probe,
transition_transaction_compaction_probe
);
let (fence_result, journal_result, decommission_target_fence_result, reconcile_result) = match fence_probe {
Ok((peer_epochs, minimum_version)) => cross_pool_fence_policy_results(peer_epochs, minimum_version),
Err(err) => {
@@ -1158,6 +1245,7 @@ pub fn start_remote_version_state_fleet_probe(topology_fingerprint: String) {
revoke_fleet_capability_proof(decommission_target_fence_fleet_proof_slot());
revoke_fleet_capability_proof(legacy_transition_state_reconcile_fleet_proof_slot());
revoke_fleet_capability_proof(ilm_recovery_export_fleet_proof_slot());
revoke_fleet_capability_proof(transition_transaction_compaction_fleet_proof_slot());
} else if let Some(err) = publish_fleet_capability_probe_result(
remote_version_state_fleet_proof_slot(),
&topology_fingerprint,
@@ -1202,6 +1290,24 @@ pub fn start_remote_version_state_fleet_probe(topology_fingerprint: String) {
"notification capability probe"
);
}
if !topology_conflict
&& let Some(err) = publish_fleet_capability_probe_result(
transition_transaction_compaction_fleet_proof_slot(),
&topology_fingerprint,
transition_transaction_compaction_result,
Instant::now(),
)
{
debug!(
event = EVENT_NOTIFICATION_CAPABILITY_PROBE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
capability = "transition_transaction_compaction_v1",
state = "failed_closed",
error = %err,
"notification capability probe"
);
}
if !topology_conflict
&& let Some(err) = publish_fleet_capability_probe_result(
tier_delete_journal_fleet_proof_slot(),
@@ -1324,6 +1430,28 @@ impl NotificationSys {
Ok(peer_epochs)
}
async fn probe_transition_transaction_compaction_fleet(&self, topology_fingerprint: &str) -> Result<BTreeMap<String, Uuid>> {
if self.peer_clients.len() != self.peer_topology_hosts.len() {
return Err(Error::other(
"transition transaction compaction capability fleet membership is incomplete",
));
}
let probes = self.peer_clients.iter().map(|client| async {
let client = client
.as_ref()
.ok_or_else(|| Error::other("transition transaction compaction capability peer is unreachable"))?;
client
.probe_transition_transaction_compaction(topology_fingerprint.to_string())
.await
});
let mut peer_epochs = BTreeMap::new();
for result in join_all(probes).await {
let (peer, epoch) = result?;
insert_remote_version_state_peer(&mut peer_epochs, peer, epoch)?;
}
Ok(peer_epochs)
}
async fn probe_cross_pool_fence_fleet(&self, topology_fingerprint: &str) -> Result<(BTreeMap<String, Uuid>, u32)> {
if self.peer_clients.len() != self.peer_topology_hosts.len() {
return Err(Error::other("cross-pool fence capability fleet membership is incomplete"));
+11
View File
@@ -883,6 +883,17 @@ pub(crate) use ops::object::body_cache_plaintext_len;
pub(crate) use ops::object::cleanup_rejected_transition_upload_durably;
#[cfg(any(test, feature = "test-util"))]
pub use ops::object::{PutObjectCommitBarrier, PutObjectCommitPause};
#[cfg(all(test, feature = "test-util"))]
pub(crate) use ops::object::{
TransitionTransactionKillPoint as SetDiskTransitionTransactionKillPoint,
TransitionTransactionKillPointBarrier as SetDiskTransitionTransactionKillPointBarrier,
};
#[cfg(all(test, feature = "test-util"))]
pub(crate) use ops::object::{
TransitionTransactionMutationKind as SetDiskTransitionTransactionMutationKind,
TransitionTransactionMutationObservation as SetDiskTransitionTransactionMutationObservation,
TransitionTransactionMutationProbe as SetDiskTransitionTransactionMutationProbe,
};
mod read;
mod replication;
pub(crate) mod shard_source;
+342 -100
View File
@@ -5666,6 +5666,10 @@ impl Drop for TransitionUploadCleanup {
if !self.armed {
return;
}
#[cfg(all(test, feature = "test-util"))]
if transition_transaction_kill_point_is_active(self.cleanup_transaction.as_ref()) {
return;
}
let Some(candidate) = self.candidate.as_ref() else {
return;
};
@@ -5870,17 +5874,29 @@ fn transition_source_identity(
}
async fn save_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction: &TransitionTransaction) -> Result<()> {
if let Some(api) = api {
return save_transition_transaction_record(api.clone(), transaction).await;
}
#[cfg(test)]
{
Ok(())
}
#[cfg(not(test))]
{
Err(Error::other("transition transaction store is unavailable"))
}
let started = std::time::Instant::now();
let result = if let Some(api) = api {
save_transition_transaction_record(api.clone(), transaction).await
} else {
#[cfg(test)]
{
Ok(())
}
#[cfg(not(test))]
{
Err(Error::other("transition transaction store is unavailable"))
}
};
#[cfg(test)]
record_transition_transaction_mutation(
transaction,
TransitionTransactionMutationKind::Create,
None,
started.elapsed(),
result.is_ok(),
);
result
}
async fn compare_and_save_transition_transaction_if_available(
@@ -5888,19 +5904,31 @@ async fn compare_and_save_transition_transaction_if_available(
expected: &TransitionTransaction,
next: &TransitionTransaction,
) -> Result<()> {
if let Some(api) = api {
#[cfg(test)]
let started = std::time::Instant::now();
let result = if let Some(api) = api {
// The transition worker already has a deep poll chain. Keep the CAS
// read/write/receipt future off Tokio's default worker stack.
return Box::pin(save_transition_transaction_record_if_current(api.clone(), expected, next)).await;
}
Box::pin(save_transition_transaction_record_if_current(api.clone(), expected, next)).await
} else {
#[cfg(test)]
{
Ok(())
}
#[cfg(not(test))]
{
Err(Error::other("transition transaction store is unavailable"))
}
};
#[cfg(test)]
{
Ok(())
}
#[cfg(not(test))]
{
Err(Error::other("transition transaction store is unavailable"))
}
record_transition_transaction_mutation(
next,
TransitionTransactionMutationKind::CompareAndSave,
Some(expected.state),
started.elapsed(),
result.is_ok(),
);
result
}
async fn advance_and_save_transition_transaction(
@@ -5909,8 +5937,6 @@ async fn advance_and_save_transition_transaction(
next: TransitionTransactionState,
remote_version: Option<TransitionRemoteVersion>,
) -> Result<()> {
#[cfg(test)]
record_transition_uploaded_save_attempt(transaction, next);
let expected = transaction.clone();
let mut advanced = expected.clone();
advanced
@@ -5922,10 +5948,33 @@ async fn advance_and_save_transition_transaction(
}
#[cfg(test)]
struct TransitionUploadedSaveProbeState {
struct TransitionTransactionMutationProbeState {
bucket: String,
object: String,
attempts: std::sync::atomic::AtomicUsize,
observations: std::sync::Mutex<Vec<TransitionTransactionMutationObservation>>,
}
#[cfg(test)]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TransitionTransactionMutationKind {
Create,
CompareAndSave,
Delete,
}
#[cfg(test)]
#[derive(Clone, Debug)]
#[allow(
dead_code,
reason = "full mutation measurements are consumed by tests behind `--features test-util`"
)]
pub(crate) struct TransitionTransactionMutationObservation {
pub(crate) kind: TransitionTransactionMutationKind,
pub(crate) previous_state: Option<TransitionTransactionState>,
pub(crate) state: TransitionTransactionState,
pub(crate) encoded_bytes: usize,
pub(crate) elapsed: std::time::Duration,
pub(crate) succeeded: bool,
}
#[cfg(test)]
@@ -5933,31 +5982,35 @@ struct TransitionUploadedSaveProbeState {
dead_code,
reason = "installed by set_disk tests behind `--features test-util` (backlog#1823)"
)]
struct TransitionUploadedSaveProbe {
state: Arc<TransitionUploadedSaveProbeState>,
pub(crate) struct TransitionTransactionMutationProbe {
state: Arc<TransitionTransactionMutationProbeState>,
}
#[cfg(test)]
static TRANSITION_UPLOADED_SAVE_PROBE: std::sync::OnceLock<std::sync::Mutex<Option<Arc<TransitionUploadedSaveProbeState>>>> =
std::sync::OnceLock::new();
static TRANSITION_TRANSACTION_MUTATION_PROBE: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<TransitionTransactionMutationProbeState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl TransitionUploadedSaveProbe {
impl TransitionTransactionMutationProbe {
#[allow(
dead_code,
reason = "installed by set_disk tests behind `--features test-util` (backlog#1823)"
)]
fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(TransitionUploadedSaveProbeState {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(TransitionTransactionMutationProbeState {
bucket: bucket.to_string(),
object: object.to_string(),
attempts: std::sync::atomic::AtomicUsize::new(0),
observations: std::sync::Mutex::new(Vec::new()),
});
let mut slot = TRANSITION_UPLOADED_SAVE_PROBE
let mut slot = TRANSITION_TRANSACTION_MUTATION_PROBE
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition uploaded-save probe mutex should not poison");
assert!(slot.is_none(), "transition uploaded-save probe must be installed by one test at a time");
.expect("transition transaction mutation probe mutex should not poison");
assert!(
slot.is_none(),
"transition transaction mutation probe must be installed by one test at a time"
);
*slot = Some(Arc::clone(&state));
drop(slot);
Self { state }
@@ -5968,17 +6021,31 @@ impl TransitionUploadedSaveProbe {
reason = "installed by set_disk tests behind `--features test-util` (backlog#1823)"
)]
fn attempts(&self) -> usize {
self.state.attempts.load(std::sync::atomic::Ordering::Acquire)
self.observations()
.into_iter()
.filter(|observation| {
observation.kind == TransitionTransactionMutationKind::CompareAndSave
&& observation.state == TransitionTransactionState::Uploaded
})
.count()
}
pub(crate) fn observations(&self) -> Vec<TransitionTransactionMutationObservation> {
self.state
.observations
.lock()
.expect("transition transaction mutation observations mutex should not poison")
.clone()
}
}
#[cfg(test)]
impl Drop for TransitionUploadedSaveProbe {
impl Drop for TransitionTransactionMutationProbe {
fn drop(&mut self) {
let mut slot = TRANSITION_UPLOADED_SAVE_PROBE
let mut slot = TRANSITION_TRANSACTION_MUTATION_PROBE
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition uploaded-save probe mutex should not poison");
.expect("transition transaction mutation probe mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
@@ -5986,19 +6053,34 @@ impl Drop for TransitionUploadedSaveProbe {
}
#[cfg(test)]
fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction, next: TransitionTransactionState) {
if next != TransitionTransactionState::Uploaded {
return;
}
let state = TRANSITION_UPLOADED_SAVE_PROBE
fn record_transition_transaction_mutation(
transaction: &TransitionTransaction,
kind: TransitionTransactionMutationKind,
previous_state: Option<TransitionTransactionState>,
elapsed: std::time::Duration,
succeeded: bool,
) {
let state = TRANSITION_TRANSACTION_MUTATION_PROBE
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition uploaded-save probe mutex should not poison")
.expect("transition transaction mutation probe mutex should not poison")
.as_ref()
.filter(|state| state.bucket == transaction.source.bucket && state.object == transaction.source.object)
.cloned();
if let Some(state) = state {
state.attempts.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
let encoded_bytes = transaction.encode().map_or(0, |encoded| encoded.len());
state
.observations
.lock()
.expect("transition transaction mutation observations mutex should not poison")
.push(TransitionTransactionMutationObservation {
kind,
previous_state,
state: transaction.state,
encoded_bytes,
elapsed,
succeeded,
});
}
}
@@ -6006,12 +6088,24 @@ async fn delete_transition_transaction_if_available(
api: Option<&Arc<ECStore>>,
transaction: &TransitionTransaction,
) -> Result<()> {
if let Some(api) = api {
#[cfg(test)]
let started = std::time::Instant::now();
let result = if let Some(api) = api {
// Conditional delete now includes a read and terminal receipt; box it
// for the same transition-worker stack bound as the CAS path above.
return Box::pin(delete_transition_transaction_record(api.clone(), transaction)).await;
}
Ok(())
Box::pin(delete_transition_transaction_record(api.clone(), transaction)).await
} else {
Ok(())
};
#[cfg(test)]
record_transition_transaction_mutation(
transaction,
TransitionTransactionMutationKind::Delete,
Some(transaction.state),
started.elapsed(),
result.is_ok(),
);
result
}
async fn delete_transition_transaction_after_remote_cleanup(
@@ -6239,6 +6333,103 @@ async fn pause_after_transition_uploaded_persisted(bucket: &str, object: &str) {
}
}
#[cfg(all(test, feature = "test-util"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TransitionTransactionKillPoint {
PrePutFence,
UploadBeforeCommitFence,
CommitFenceBeforeLocalCommit,
LocalCommitBeforeDelete,
}
#[cfg(all(test, feature = "test-util"))]
struct TransitionTransactionKillPointBarrierState {
bucket: String,
object: String,
point: TransitionTransactionKillPoint,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(all(test, feature = "test-util"))]
pub(crate) struct TransitionTransactionKillPointBarrier {
state: Arc<TransitionTransactionKillPointBarrierState>,
}
#[cfg(all(test, feature = "test-util"))]
static TRANSITION_TRANSACTION_KILL_POINT_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<TransitionTransactionKillPointBarrierState>>>,
> = std::sync::OnceLock::new();
#[cfg(all(test, feature = "test-util"))]
impl TransitionTransactionKillPointBarrier {
pub(crate) fn install(bucket: &str, object: &str, point: TransitionTransactionKillPoint) -> Self {
let state = Arc::new(TransitionTransactionKillPointBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
point,
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = TRANSITION_TRANSACTION_KILL_POINT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition transaction kill-point barrier mutex should not poison");
assert!(slot.is_none(), "one transition transaction kill-point may be installed at a time");
*slot = Some(Arc::clone(&state));
drop(slot);
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("transition should reach the requested transaction kill-point");
}
}
#[cfg(all(test, feature = "test-util"))]
impl Drop for TransitionTransactionKillPointBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = TRANSITION_TRANSACTION_KILL_POINT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition transaction kill-point barrier mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(all(test, feature = "test-util"))]
async fn pause_transition_transaction_at(bucket: &str, object: &str, point: TransitionTransactionKillPoint) {
let barrier = TRANSITION_TRANSACTION_KILL_POINT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition transaction kill-point barrier mutex should not poison")
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object && barrier.point == point)
.cloned();
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
}
}
#[cfg(all(test, feature = "test-util"))]
fn transition_transaction_kill_point_is_active(transaction: Option<&TransitionTransaction>) -> bool {
let Some(transaction) = transaction else {
return false;
};
TRANSITION_TRANSACTION_KILL_POINT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("transition transaction kill-point barrier mutex should not poison")
.as_ref()
.is_some_and(|barrier| barrier.bucket == transaction.source.bucket && barrier.object == transaction.source.object)
}
#[cfg(test)]
#[derive(Clone, Copy, PartialEq, Eq)]
enum TransitionCommitPause {
@@ -8987,7 +9178,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let oi = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended);
let transaction_api = transition_object_store(&self.ctx).await;
let mut transaction = TransitionTransaction::new(TransitionTransactionInit {
let transition_compaction_fleet_proof =
crate::services::notification_sys::acquire_transition_transaction_compaction_fleet_proof();
let compact_transition_transaction = transition_compaction_fleet_proof
.as_ref()
.is_some_and(crate::services::notification_sys::transition_transaction_compaction_fleet_proof_matches);
let transaction_init = TransitionTransactionInit {
deployment_id: transition_deployment_id(&self.ctx)?,
transaction_id: Uuid::new_v4(),
owner_epoch: Uuid::new_v4(),
@@ -8996,9 +9192,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
tier_name: opts.transition.tier.clone(),
backend_fingerprint: tgt_client.backend_identity(),
not_after_unix_nanos: transition_transaction_not_after_unix_nanos()?,
})
};
let mut transaction = if compact_transition_transaction {
TransitionTransaction::new_compact(transaction_init)
} else {
TransitionTransaction::new(transaction_init)
}
.map_err(Error::other)?;
save_transition_transaction_if_available(transaction_api.as_ref(), &transaction).await?;
#[cfg(all(test, feature = "test-util"))]
pause_transition_transaction_at(bucket, object, TransitionTransactionKillPoint::PrePutFence).await;
let transaction_id = transaction.transaction_id;
let dest_obj = transaction.remote_object.clone();
let mut transition_meta = (*oi.user_defined).clone();
@@ -9075,14 +9278,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let mut upload_cleanup = TransitionUploadCleanup::new(tgt_client, &dest_obj);
upload_cleanup.set_cleanup_owner(transaction_api.clone(), &transaction);
advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::UploadOutcomeUnknown,
None,
)
.await?;
upload_cleanup.update_cleanup_transaction(&transaction);
if !compact_transition_transaction {
advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::UploadOutcomeUnknown,
None,
)
.await?;
upload_cleanup.update_cleanup_transaction(&transaction);
}
let remote_upload = {
let lease = &upload_cleanup.lease;
let recorded_candidate = &mut upload_cleanup.candidate;
@@ -9142,27 +9347,31 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
return Err(err.into());
}
};
if let Err(err) = advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::Uploaded,
Some(TransitionRemoteVersion::known_from_put_response(candidate.remote_version().to_string())),
)
.await
{
let cleanup_api = transition_cleanup_store(&self.ctx).await;
if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api, &mut transaction).await {
return Err(StorageError::Io(std::io::Error::other(format!(
"{err}; uploaded transition transaction persist failed and cleanup failed: {cleanup_err}"
))));
if !compact_transition_transaction {
if let Err(err) = advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::Uploaded,
Some(TransitionRemoteVersion::known_from_put_response(candidate.remote_version().to_string())),
)
.await
{
let cleanup_api = transition_cleanup_store(&self.ctx).await;
if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api, &mut transaction).await {
return Err(StorageError::Io(std::io::Error::other(format!(
"{err}; uploaded transition transaction persist failed and cleanup failed: {cleanup_err}"
))));
}
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), &transaction, bucket, object).await;
return Err(err);
}
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), &transaction, bucket, object).await;
return Err(err);
upload_cleanup.update_cleanup_transaction(&transaction);
}
upload_cleanup.update_cleanup_transaction(&transaction);
#[cfg(all(test, feature = "test-util"))]
pause_after_transition_uploaded_persisted(bucket, object).await;
#[cfg(all(test, feature = "test-util"))]
pause_transition_transaction_at(bucket, object, TransitionTransactionKillPoint::UploadBeforeCommitFence).await;
let commit_opts = opts.as_commit_opts();
// Note: Using clone() here is necessary because ObjectOptions has 124 fields.
@@ -9273,11 +9482,25 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
return Err(Error::other("remote version state fleet capability changed during transition"));
}
if compact_transition_transaction
&& !transition_compaction_fleet_proof
.as_ref()
.is_some_and(crate::services::notification_sys::transition_transaction_compaction_fleet_proof_matches)
{
drop(transition_lock_guard);
if upload_cleanup.cleanup().await.is_ok() {
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), &transaction, bucket, object).await;
}
return Err(Error::other(
"transition transaction compaction fleet capability changed during transition",
));
}
if let Err(err) = advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::LocalCommitStarted,
None,
compact_transition_transaction
.then(|| TransitionRemoteVersion::known_from_put_response(candidate.remote_version().to_string())),
)
.await
{
@@ -9287,6 +9510,9 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
return Err(err);
}
upload_cleanup.update_cleanup_transaction(&transaction);
#[cfg(all(test, feature = "test-util"))]
pause_transition_transaction_at(bucket, object, TransitionTransactionKillPoint::CommitFenceBeforeLocalCommit).await;
upload_cleanup.disarm();
if let Err(err) = self.delete_object_version(bucket, object, &fi, false).await {
warn!(
@@ -9299,34 +9525,48 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
drop(transition_lock_guard);
return Err(err);
}
match advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::Committed,
None,
)
.await
{
Ok(()) => {
if let Err(err) = delete_transition_transaction_if_available(transaction_api.as_ref(), &transaction).await {
warn!(
bucket = bucket,
object = object,
transaction_id = %transaction_id,
error = ?err,
"transition committed locally but transaction cleanup failed"
);
}
}
Err(err) => {
#[cfg(all(test, feature = "test-util"))]
pause_transition_transaction_at(bucket, object, TransitionTransactionKillPoint::LocalCommitBeforeDelete).await;
if compact_transition_transaction {
if let Err(err) = delete_transition_transaction_if_available(transaction_api.as_ref(), &transaction).await {
warn!(
bucket = bucket,
object = object,
transaction_id = %transaction_id,
error = ?err,
"transition committed locally but transaction committed-state advance failed"
"transition committed locally but compact transaction cleanup failed"
);
}
} else {
match advance_and_save_transition_transaction(
transaction_api.as_ref(),
&mut transaction,
TransitionTransactionState::Committed,
None,
)
.await
{
Ok(()) => {
if let Err(err) = delete_transition_transaction_if_available(transaction_api.as_ref(), &transaction).await {
warn!(
bucket = bucket,
object = object,
transaction_id = %transaction_id,
error = ?err,
"transition committed locally but transaction cleanup failed"
);
}
}
Err(err) => {
warn!(
bucket = bucket,
object = object,
transaction_id = %transaction_id,
error = ?err,
"transition committed locally but transaction committed-state advance failed"
);
}
}
}
// delete_object_version persisted transition_status=complete and freed the
@@ -14929,6 +15169,7 @@ mod transition_upload_integrity_tests {
use super::*;
use crate::bucket::lifecycle::lifecycle::{TRANSITION_PENDING, TransitionOptions};
use crate::layout::endpoints::SetupType;
use crate::services::notification_sys::install_transition_transaction_compaction_fleet_proof_for_test;
use crate::services::tier::test_util::register_mock_tier;
use crate::set_disk::replication::RestoreFinalizeBarrier;
use http::HeaderMap;
@@ -16195,6 +16436,7 @@ mod transition_upload_integrity_tests {
let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
let _compaction_proof = install_transition_transaction_compaction_fleet_proof_for_test("object-transaction-fencing-test");
let barrier = TransitionCommitBarrier::install(bucket, object);
let transition_set = Arc::clone(&set_disks);
@@ -16397,7 +16639,7 @@ mod transition_upload_integrity_tests {
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
backend.set_put_remote_version(Some(String::new())).await;
let save_probe = TransitionUploadedSaveProbe::install(bucket, object);
let save_probe = TransitionTransactionMutationProbe::install(bucket, object);
set_disks
.transition_object(bucket, object, &transition_options(&original, tier_name))
@@ -16461,7 +16703,7 @@ mod transition_upload_integrity_tests {
let remote_version = Uuid::nil().to_string();
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
backend.set_put_remote_version(Some(remote_version.clone())).await;
let save_probe = TransitionUploadedSaveProbe::install(bucket, object);
let save_probe = TransitionTransactionMutationProbe::install(bucket, object);
set_disks
.transition_object(bucket, object, &transition_options(&original, tier_name))
+369 -2
View File
@@ -877,7 +877,11 @@ mod tests {
data_movement::SourceCleanupDeleteBarrier,
disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE},
runtime::{global::set_object_store_resolver, sources as runtime_sources},
services::notification_sys::acquire_tier_delete_journal_fleet_proof,
services::notification_sys::{
acquire_tier_delete_journal_fleet_proof, acquire_transition_transaction_compaction_fleet_proof,
install_transition_transaction_compaction_fleet_proof_for_test,
transition_transaction_compaction_fleet_proof_matches,
},
services::tier::{
test_util::{MockWarmBackend, MockWarmOp, TransitionCleanupStoreBarrier, register_mock_tier},
tier::{
@@ -899,7 +903,14 @@ mod tests {
},
warm_backend::{TransitionCandidateProbe, WarmBackend},
},
set_disk::SetDiskTransitionUploadedCommitBarrier as TransitionUploadedCommitBarrier,
set_disk::{
SetDiskTransitionTransactionKillPoint as TransitionTransactionKillPoint,
SetDiskTransitionTransactionKillPointBarrier as TransitionTransactionKillPointBarrier,
SetDiskTransitionTransactionMutationKind as TransitionTransactionMutationKind,
SetDiskTransitionTransactionMutationObservation as TransitionTransactionMutationObservation,
SetDiskTransitionTransactionMutationProbe as TransitionTransactionMutationProbe,
SetDiskTransitionUploadedCommitBarrier as TransitionUploadedCommitBarrier,
},
storage_api_contracts::list::ListOperations as _,
};
#[cfg(feature = "test-util")]
@@ -11586,6 +11597,36 @@ mod tests {
.len()
}
#[cfg(feature = "test-util")]
async fn only_transition_transaction(store: Arc<crate::store::ECStore>) -> TransitionTransaction {
let records = store
.clone()
.list_objects_v2(
RUSTFS_META_BUCKET,
TRANSITION_TRANSACTION_RECORD_PREFIX,
None,
None,
100,
false,
None,
false,
)
.await
.expect("transition transaction records should be listable")
.objects;
assert_eq!(records.len(), 1, "test fixture should have exactly one transition transaction");
let transaction_id = records[0]
.name
.rsplit('/')
.next()
.and_then(|name| name.strip_suffix(".json"))
.and_then(|name| uuid::Uuid::parse_str(name).ok())
.expect("transition transaction path should end in its UUID");
load_transition_transaction_record(store, transaction_id)
.await
.expect("transition transaction should load")
}
#[cfg(feature = "test-util")]
async fn register_transition_reconcile_test_tier(
handle: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
@@ -19585,6 +19626,330 @@ mod tests {
assert_eq!(loaded_ids, intent_ids);
}
#[cfg(feature = "test-util")]
fn transition_mutation_measurement(observations: &[TransitionTransactionMutationObservation]) -> (usize, usize, u128) {
(
observations.len(),
observations.iter().map(|observation| observation.encoded_bytes).sum(),
observations.iter().map(|observation| observation.elapsed.as_micros()).sum(),
)
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn compact_transition_transactions_halve_success_path_quorum_mutations() {
let temp_dir = tempfile::tempdir().expect("create transition mutation measurement dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "transition-mutation-measurement", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "MUTATION-MEASURE";
register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let bucket = "transition-mutation-measurement";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("measurement bucket should be created");
let run_case = |profile: &'static str, size: usize| {
let store = store.clone();
async move {
let object = format!("{profile}-{size}.bin");
let mut reader = PutObjReader::from_vec(vec![b'm'; size]);
let original = store
.put_object(bucket, &object, &mut reader, &ObjectOptions::default())
.await
.expect("measurement source should be written");
let probe = TransitionTransactionMutationProbe::install(bucket, &object);
store
.transition_object(
bucket,
&object,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: original.etag.clone().expect("measurement source should have an ETag"),
..Default::default()
},
version_id: original.version_id.map(|version| version.to_string()),
mod_time: original.mod_time,
..Default::default()
},
)
.await
.expect("measurement transition should commit");
probe.observations()
}
};
let sizes = [4 * 1024, 1024 * 1024];
let mut legacy = Vec::with_capacity(sizes.len());
for size in sizes {
legacy.push((size, run_case("legacy", size).await));
}
let compaction_proof = install_transition_transaction_compaction_fleet_proof_for_test("object-transaction-fencing-test");
for ((size, legacy_observations), compact_size) in legacy.into_iter().zip(sizes) {
assert_eq!(size, compact_size);
let compact_observations = run_case("compact", size).await;
let legacy_measurement = transition_mutation_measurement(&legacy_observations);
let compact_measurement = transition_mutation_measurement(&compact_observations);
assert_eq!(legacy_measurement.0, 6, "legacy success should use five saves and one delete");
assert_eq!(compact_measurement.0, 3, "compact success should use two saves and one delete");
assert!(
compact_measurement.1 < legacy_measurement.1,
"compact transaction bodies should write fewer aggregate bytes"
);
assert_eq!(
legacy_observations
.iter()
.map(|observation| (observation.kind, observation.previous_state, observation.state))
.collect::<Vec<_>>(),
vec![
(TransitionTransactionMutationKind::Create, None, TransitionTransactionState::UploadStarted),
(
TransitionTransactionMutationKind::CompareAndSave,
Some(TransitionTransactionState::UploadStarted),
TransitionTransactionState::UploadOutcomeUnknown,
),
(
TransitionTransactionMutationKind::CompareAndSave,
Some(TransitionTransactionState::UploadOutcomeUnknown),
TransitionTransactionState::Uploaded,
),
(
TransitionTransactionMutationKind::CompareAndSave,
Some(TransitionTransactionState::Uploaded),
TransitionTransactionState::LocalCommitStarted,
),
(
TransitionTransactionMutationKind::CompareAndSave,
Some(TransitionTransactionState::LocalCommitStarted),
TransitionTransactionState::Committed,
),
(
TransitionTransactionMutationKind::Delete,
Some(TransitionTransactionState::Committed),
TransitionTransactionState::Committed,
),
]
);
assert_eq!(
compact_observations
.iter()
.map(|observation| (observation.kind, observation.previous_state, observation.state))
.collect::<Vec<_>>(),
vec![
(
TransitionTransactionMutationKind::Create,
None,
TransitionTransactionState::UploadOutcomeUnknown,
),
(
TransitionTransactionMutationKind::CompareAndSave,
Some(TransitionTransactionState::UploadOutcomeUnknown),
TransitionTransactionState::LocalCommitStarted,
),
(
TransitionTransactionMutationKind::Delete,
Some(TransitionTransactionState::LocalCommitStarted),
TransitionTransactionState::LocalCommitStarted,
),
]
);
assert!(legacy_observations.iter().all(|observation| observation.succeeded));
assert!(compact_observations.iter().all(|observation| observation.succeeded));
println!(
"transition_mutation_measurement,profile=legacy,size={size},mutations={},encoded_bytes={},latency_us={},quorum_operations={}",
legacy_measurement.0, legacy_measurement.1, legacy_measurement.2, legacy_measurement.0
);
println!(
"transition_mutation_measurement,profile=compact,size={size},mutations={},encoded_bytes={},latency_us={},quorum_operations={}",
compact_measurement.0, compact_measurement.1, compact_measurement.2, compact_measurement.0
);
}
assert_eq!(transition_transaction_record_count(store.clone()).await, 0);
let admitted = acquire_transition_transaction_compaction_fleet_proof()
.expect("published homogeneous proof should admit one compact writer");
assert!(transition_transaction_compaction_fleet_proof_matches(&admitted));
drop(compaction_proof);
assert!(
!transition_transaction_compaction_fleet_proof_matches(&admitted),
"revocation must fence a writer admitted by the previous process-epoch snapshot"
);
drop(admitted);
assert!(
acquire_transition_transaction_compaction_fleet_proof().is_none(),
"revoking the homogeneous proof must restore the legacy writer profile"
);
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn compact_transition_kill_points_preserve_the_only_remote_owner() {
let temp_dir = tempfile::tempdir().expect("create compact transition kill-point dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "compact-transition-kill-points", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "COMPACT-KILL";
let backend = register_transition_reconcile_test_tier(&ctx.tier_config_mgr(), tier_name).await;
let _tier_lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("mock tier lease should remain available during recovery");
let bucket = "compact-transition-kill-points";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("kill-point bucket should be created");
let _compaction_proof = install_transition_transaction_compaction_fleet_proof_for_test("object-transaction-fencing-test");
for (index, point, expected_state, expected_revision, expected_committed, expected_recovered) in [
(
0,
TransitionTransactionKillPoint::PrePutFence,
TransitionTransactionState::UploadOutcomeUnknown,
1,
false,
true,
),
(
1,
TransitionTransactionKillPoint::UploadBeforeCommitFence,
TransitionTransactionState::UploadOutcomeUnknown,
1,
false,
true,
),
(
2,
TransitionTransactionKillPoint::LocalCommitBeforeDelete,
TransitionTransactionState::LocalCommitStarted,
2,
true,
true,
),
(
3,
TransitionTransactionKillPoint::CommitFenceBeforeLocalCommit,
TransitionTransactionState::LocalCommitStarted,
2,
false,
false,
),
] {
let object = format!("kill-point-{index}.bin");
let payload = vec![b'k' + u8::try_from(index).expect("small case index should fit u8"); 64 * 1024];
let mut reader = PutObjReader::from_vec(payload.clone());
let source = store
.put_object(bucket, &object, &mut reader, &ObjectOptions::default())
.await
.expect("kill-point source should be written");
let remote_before = backend.object_count().await;
let barrier = TransitionTransactionKillPointBarrier::install(bucket, &object, point);
let transition_store = store.clone();
let transition_object = object.clone();
let transition = tokio::spawn(async move {
transition_store
.transition_object(
bucket,
&transition_object,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: source.etag.clone().expect("kill-point source should have an ETag"),
..Default::default()
},
version_id: source.version_id.map(|version| version.to_string()),
mod_time: source.mod_time,
..Default::default()
},
)
.await
});
barrier.wait_until_paused().await;
transition.abort();
assert!(
transition
.await
.expect_err("kill-point transition should be cancelled")
.is_cancelled(),
"kill-point transition should stop without unwinding"
);
drop(barrier);
let transaction = only_transition_transaction(store.clone()).await;
assert_eq!((transaction.state, transaction.revision), (expected_state, expected_revision));
let paused_source = store
.get_object_info(
bucket,
&object,
&ObjectOptions {
metadata_cache_safe: false,
..Default::default()
},
)
.await
.expect("kill-point source metadata should remain readable");
assert_eq!(
paused_source.transitioned_object.status == rustfs_filemeta::TRANSITION_COMPLETE,
expected_committed
);
let expected_remote_at_pause = remote_before + usize::from(point != TransitionTransactionKillPoint::PrePutFence);
assert_eq!(backend.object_count().await, expected_remote_at_pause);
let stats = recover_transition_transaction_records_at(
store.clone(),
100,
None,
i128::from(transaction.not_after_unix_nanos) + 1,
)
.await
.expect("kill-point transaction recovery should complete");
if expected_recovered {
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 1, 0, 0));
assert_eq!(transition_transaction_record_count(store.clone()).await, 0);
} else {
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 0, 1, 0));
assert_eq!(
transition_transaction_record_count(store.clone()).await,
1,
"an uncommitted local-commit fence must retain its exact remote owner"
);
}
let expected_remote_after_recovery = if matches!(
point,
TransitionTransactionKillPoint::LocalCommitBeforeDelete
| TransitionTransactionKillPoint::CommitFenceBeforeLocalCommit
) {
remote_before + 1
} else {
remote_before
};
assert_eq!(backend.object_count().await, expected_remote_after_recovery);
assert_eq!(backend.remove_count().await, usize::from(index >= 1));
let mut restored = Vec::new();
store
.get_object_reader(bucket, &object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("kill-point source should remain readable through its authoritative location")
.stream
.read_to_end(&mut restored)
.await
.expect("kill-point source body should drain");
assert_eq!(restored, payload);
if !expected_recovered {
break;
}
}
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
@@ -21236,6 +21601,7 @@ mod tests {
let tier_name = "TXRESPONSELOSS";
let backend = register_transition_reconcile_test_tier(&ctx.tier_config_mgr(), tier_name).await;
let _compaction_proof = install_transition_transaction_compaction_fleet_proof_for_test("object-transaction-fencing-test");
let bucket = "transition-response-loss-bucket";
let object = "source.bin";
store
@@ -21304,6 +21670,7 @@ mod tests {
TransitionTransactionState::UploadOutcomeUnknown,
"a response-lost PUT must not remain in UploadStarted"
);
assert_eq!(transaction.revision, 1, "compact response loss must retain the pre-PUT fence generation");
assert!(
backend.contains(&transaction.remote_object).await,
"the test backend must retain the remote candidate"
+37 -6
View File
@@ -176,6 +176,8 @@ pub const HEAL_CONTROL_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-heal-control-ca
pub const REMOTE_VERSION_STATE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-tier-remote-version-state-capability-v1\0";
pub const CROSS_POOL_FENCE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-cross-pool-fence-capability-v1\0";
pub const ILM_RECOVERY_EXPORT_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-ilm-recovery-export-capability-v1\0";
pub const TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX: &[u8] =
b"rustfs-transition-transaction-compaction-capability-v1\0";
pub const TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE: usize = 64 * 1024;
pub const TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE: usize = 1024;
pub const TIER_MUTATION_RPC_MAX_ABORT_PAYLOAD_SIZE: usize = TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE;
@@ -232,6 +234,18 @@ pub fn is_ilm_recovery_export_capability_probe(command: &[u8]) -> bool {
&& command.starts_with(ILM_RECOVERY_EXPORT_CAPABILITY_PROBE_PREFIX)
}
pub fn transition_transaction_compaction_capability_probe(nonce: &[u8; 16]) -> Vec<u8> {
let mut probe = Vec::with_capacity(TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX.len() + nonce.len());
probe.extend_from_slice(TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX);
probe.extend_from_slice(nonce);
probe
}
pub fn is_transition_transaction_compaction_capability_probe(command: &[u8]) -> bool {
command.len() == TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX.len() + 16
&& command.starts_with(TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX)
}
pub fn encode_remote_version_state_capability(
topology_member: &str,
process_epoch: &[u8; 16],
@@ -2152,12 +2166,14 @@ mod heal_control_tests {
use super::{
CROSS_POOL_FENCE_CAPABILITY_PROBE_PREFIX, HEAL_CONTROL_CAPABILITY_PROBE_PREFIX, HEAL_CONTROL_PROTOCOL_VERSION,
ILM_RECOVERY_EXPORT_CAPABILITY_PROBE_PREFIX, REMOTE_VERSION_STATE_CAPABILITY_PROBE_PREFIX,
canonical_heal_control_capability_ack, canonical_heal_control_request_body, canonical_heal_control_response_body,
decode_remote_version_state_capability, encode_cross_pool_fence_capability, encode_remote_version_state_capability,
heal_control_capability_probe, heal_control_coordinator_epoch, heal_control_execution_timeout,
heal_control_execution_timeout_for, ilm_recovery_export_capability_probe, internode_rpc_timeout,
is_cross_pool_fence_capability_probe, is_heal_control_capability_probe, is_ilm_recovery_export_capability_probe,
is_remote_version_state_capability_probe, normalize_internode_rpc_timeout, remote_version_state_capability_probe,
TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX, canonical_heal_control_capability_ack,
canonical_heal_control_request_body, canonical_heal_control_response_body, decode_remote_version_state_capability,
encode_cross_pool_fence_capability, encode_remote_version_state_capability, heal_control_capability_probe,
heal_control_coordinator_epoch, heal_control_execution_timeout, heal_control_execution_timeout_for,
ilm_recovery_export_capability_probe, internode_rpc_timeout, is_cross_pool_fence_capability_probe,
is_heal_control_capability_probe, is_ilm_recovery_export_capability_probe, is_remote_version_state_capability_probe,
is_transition_transaction_compaction_capability_probe, normalize_internode_rpc_timeout,
remote_version_state_capability_probe, transition_transaction_compaction_capability_probe,
};
use crate::heal_control;
use std::time::Duration;
@@ -2234,6 +2250,21 @@ mod heal_control_tests {
assert!(!is_ilm_recovery_export_capability_probe(&extra));
}
#[test]
fn transition_transaction_compaction_probe_requires_exact_prefix_and_nonce() {
let probe = transition_transaction_compaction_capability_probe(&[7; 16]);
assert!(is_transition_transaction_compaction_capability_probe(&probe));
assert!(!is_transition_transaction_compaction_capability_probe(
TRANSITION_TRANSACTION_COMPACTION_CAPABILITY_PROBE_PREFIX,
));
let mut wrong_prefix = probe.clone();
wrong_prefix[0] ^= 1;
assert!(!is_transition_transaction_compaction_capability_probe(&wrong_prefix));
let mut extra = probe;
extra.push(0);
assert!(!is_transition_transaction_compaction_capability_probe(&extra));
}
#[test]
fn remote_version_state_capability_binds_member_and_process_epoch() {
let encoded =