Compare commits

..

2 Commits

Author SHA1 Message Date
Zhengchao An 2cbb5937d0 Merge branch 'main' into cxymds/perf/ilm-transition-mutation-reduction 2026-09-07 02:11:30 +08:00
cxymds dac6915c08 perf(ilm): reduce transition transaction mutations 2026-09-07 00:00:40 +08:00
13 changed files with 1135 additions and 200 deletions
@@ -3837,77 +3837,6 @@ async fn test_bucket_replication_converges_delete_marker_and_version_purge() ->
Ok(())
}
/// Regression for rustfs/backlog#2340 (not Wasabi specific): a directory
/// marker (`prefix/` with a body) in a versioned bucket is stored as the null
/// version, like MinIO (`putOpts`: "for directory objects skip creating new
/// versions"), and must still replicate to completion instead of staying
/// `PENDING`.
#[tokio::test]
async fn test_bucket_replication_replicates_directory_marker_in_versioned_bucket() -> TestResult {
init_logging();
let mut source_env = RustFSTestEnvironment::new().await?;
let mut source_env_vars = replication_fast_env();
source_env_vars.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
source_env.start_rustfs_server_with_env(vec![], &source_env_vars).await?;
let mut target_env = RustFSTestEnvironment::new().await?;
target_env.start_rustfs_server_without_cleanup(vec![]).await?;
let source_bucket = "replication-dir-marker-src";
let target_bucket = "replication-dir-marker-dst";
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
target_client.create_bucket().bucket(target_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
enable_bucket_versioning(&target_env, target_bucket).await?;
let target_arn = set_replication_target(&source_env, source_bucket, &target_env, target_bucket).await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let marker_key = "dir/trailing/";
let body = b"directory marker body";
let put = source_client
.put_object()
.bucket(source_bucket)
.key(marker_key)
.body(ByteStream::from_static(body))
.send()
.await?;
assert!(
put.version_id()
.is_none_or(|id| id == "null" || id == uuid::Uuid::nil().to_string()),
"a directory marker is the null version even in a versioned bucket: {:?}",
put.version_id()
);
wait_for_source_replication_status(&source_client, source_bucket, marker_key, "COMPLETED", false).await?;
let replica = target_client
.get_object()
.bucket(target_bucket)
.key(marker_key)
.send()
.await?;
assert_eq!(replica.body.collect().await?.into_bytes().as_ref(), body);
let listed = target_client
.list_object_versions()
.bucket(target_bucket)
.prefix(marker_key)
.send()
.await?;
let marker_versions: Vec<_> = listed.versions().iter().filter(|v| v.key() == Some(marker_key)).collect();
assert_eq!(marker_versions.len(), 1, "the marker must land exactly once: {marker_versions:?}");
assert_eq!(
marker_versions[0].version_id(),
Some("null"),
"the replica keeps the null version identity"
);
Ok(())
}
#[tokio::test]
async fn test_bucket_replication_disabled_delete_marker_does_not_propagate() -> TestResult {
init_logging();
@@ -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),
(
@@ -1954,7 +1970,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)
@@ -1966,7 +1982,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 {
@@ -2384,6 +2402,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();
@@ -1697,6 +1697,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
@@ -317,12 +317,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);
@@ -350,6 +359,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();
@@ -450,6 +463,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()
@@ -1072,6 +1112,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"));
@@ -1089,6 +1158,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);
}
@@ -1132,8 +1202,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) => {
@@ -1157,6 +1244,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,
@@ -1201,6 +1289,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(),
@@ -1323,6 +1429,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 {
AfterPrePutFence,
AfterUploadBeforeCommitFence,
AfterCommitFenceBeforeLocalCommit,
AfterLocalCommitBeforeDelete,
}
#[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::AfterPrePutFence).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::AfterUploadBeforeCommitFence).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::AfterCommitFenceBeforeLocalCommit).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::AfterLocalCommitBeforeDelete).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
@@ -870,7 +870,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::{
@@ -892,7 +896,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")]
@@ -11579,6 +11590,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>>,
@@ -19283,6 +19324,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::AfterPrePutFence,
TransitionTransactionState::UploadOutcomeUnknown,
1,
false,
true,
),
(
1,
TransitionTransactionKillPoint::AfterUploadBeforeCommitFence,
TransitionTransactionState::UploadOutcomeUnknown,
1,
false,
true,
),
(
2,
TransitionTransactionKillPoint::AfterLocalCommitBeforeDelete,
TransitionTransactionState::LocalCommitStarted,
2,
true,
true,
),
(
3,
TransitionTransactionKillPoint::AfterCommitFenceBeforeLocalCommit,
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::AfterPrePutFence);
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::AfterLocalCommitBeforeDelete
| TransitionTransactionKillPoint::AfterCommitFenceBeforeLocalCommit
) {
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)]
@@ -20867,6 +21232,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
@@ -20935,6 +21301,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],
@@ -2141,12 +2155,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;
@@ -2223,6 +2239,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 =
@@ -41,7 +41,7 @@ All keys below are objects in the internal metadata bucket. The table gives the
| Protocol | Current schema/version | Canonical key | Creator and cleanup owner | Authoritative identity and mutable fields | Current durability point |
|---|---|---|---|---|---|
| Transition transaction | `rustfs-transition-transaction-v1`; successor v2 is approved below but not implemented | `ilm/transition-transactions/records/<aa>/<bb>/<transaction-id>.json` | The transition attempt creates it; transition commit/recovery cleans it | Immutable/fence identity: deployment, transaction, fixed v1 `owner_epoch`, write, source identity, tier/backend fingerprint, canonical remote object, deadline. Mutable: state, remote version, revision. `TransitionCleanupProof` is only a transient admission input to `mark_cleanup_pending`; it is not persisted in the record | Create-only maximum-parity write; exact record and ETag read before successor `If-Match`; terminal receipt followed by exact ETag conditional delete |
| Transition transaction | `rustfs-transition-transaction-v1`; the compact v1 state profile is fleet-gated; successor v2 is approved below but not implemented | `ilm/transition-transactions/records/<aa>/<bb>/<transaction-id>.json` | The transition attempt creates it; transition commit/recovery cleans it | Immutable/fence identity: deployment, transaction, fixed v1 `owner_epoch`, write, source identity, tier/backend fingerprint, canonical remote object, deadline. Mutable: state, remote version, revision. `TransitionCleanupProof` is only a transient admission input to `mark_cleanup_pending`; it is not persisted in the record | Create-only maximum-parity write; exact record and ETag read before successor `If-Match`; terminal receipt followed by exact ETag conditional delete. A compact success uses two saves and one delete instead of five saves and one delete |
| Tier mutation peer intent | `rustfs-tier-mutation-intent-v1` | `tier/mutation-intents/records/<aa>/<bb>/<mutation-id>.json` | The receiving peer creates and converges it; the mutation recovery path cleans it | Immutable: mutation ID/kind, old config ETag, candidate digest, sorted affected target identities, expiry. Mutable: revision, state, committed config ETag | Create with `If-None-Match: *`; transition/delete with ETag `If-Match`; maximum parity |
| Tier mutation coordinator intent | `rustfs-tier-mutation-intent-v1` | `tier/mutation-intents/coordinators/<aa>/<bb>/<mutation-id>.json` | The initiating node creates it; coordinator recovery cleans it after peer convergence | Same mutation identity and mutable fields as the peer record | Same conditional-write contract as the peer intent |
| Tier validation probe intent | Dormant `rustfs-tier-probe-intent-v1`; no writer or recovery is enabled | `ilm/tier-probe-intents/records/<aa>/<bb>/<probe-id>.json` | No current runtime owner because no path creates the record; v1 permits only the immutable creator as owner | Immutable probe, operation-generation, destination, random remote object, creator identity, and v1 owner fence. Mutable: revision, state, and monotonic remote-version proof | Conditional create/CAS/delete primitives exist but are not called by Add/Edit/Verify or recovery |
@@ -109,7 +109,7 @@ Tier mutation backend validation is outside both exclusive guards and is bound t
### Current contract
`TransitionTransaction` binds the canonical candidate name to a transaction UUID, write UUID, source identity, tier name, backend fingerprint, remote-version state, deadline, fixed `owner_epoch`, and mutable `revision`. The state model permits these ordinary edges:
`TransitionTransaction` binds the canonical candidate name to a transaction UUID, write UUID, source identity, tier name, backend fingerprint, remote-version state, deadline, fixed `owner_epoch`, and mutable `revision`. Without a current homogeneous compaction capability proof, writers use the legacy state profile:
```text
UploadStarted -> Uploaded -> LocalCommitStarted -> Committed
@@ -117,6 +117,16 @@ UploadStarted -> Uploaded -> LocalCommitStarted -> Committed
\-> AbortedNoRemote
```
When every current topology member answers the exact `transition_transaction_compaction_v1` capability challenge, a writer may hold that generation's non-cloneable proof permit and emit the compact v1 profile:
```text
UploadOutcomeUnknown@1 -> LocalCommitStarted@2 -> conditional record delete
```
The create-only `UploadOutcomeUnknown@1` record is durable before the remote PUT. After PUT succeeds, the writer revalidates the source, Object Lock decision, tier generation, and fleet proof before one exact successor CAS to `LocalCommitStarted@2`; that successor carries the known remote version. After the local metadata commit, the writer conditionally deletes that exact record rather than persisting a redundant `Committed` generation. A crash before the PUT leaves an unknown record whose provider probe can prove absence. A crash after PUT leaves the canonical candidate under the unknown record. A crash after the commit fence leaves the exact remote tuple under `LocalCommitStarted`, and a crash after local commit lets recovery prove ownership transfer and remove only the record.
The `UploadOutcomeUnknown@1 -> LocalCommitStarted@2` pair is the only compact direct edge. Legacy `UploadOutcomeUnknown@2 -> LocalCommitStarted@3` remains invalid, while historical receipt jumps keep their separately enumerated revision distances. This distinction lets existing v1 payload readers decode compact records without treating arbitrary skipped history as valid. If any peer is absent, old, unreachable, restarted with a different process epoch, or the topology changes, proof publication/revalidation fails closed and new attempts use the legacy profile. Planned downgrade first disables compact admission and drains live compact records; an older runtime reader retains or safely reconciles the v1 states, but an older decommission checkpoint validator can reject the compact successor and block pool completion.
Separately, `mark_cleanup_pending` permits proof-checked model edges from `Uploaded`, `UploadOutcomeUnknown`, and `LocalCommitStarted`. Current production recovery emits `CleanupPending` after an expired `Uploaded` record wins the exact successor CAS, or when an expired `UploadOutcomeUnknown` probe returns `UnversionedPresent` or `VersionedPresent` with a non-nil identifier. `LocalCommitStarted` mismatch or missing-source recovery retains the record; that cleanup edge is currently exercised through the state-machine API and tests, not produced by runtime recovery. States that require a remote delete still require a known `TransitionRemoteVersion` kind. A probed versioned candidate whose identifier parses as a nil UUID is retained and never authorizes remote deletion.
The remote candidate itself is named by `canonical_transition_remote_object` under `ilm/transition-transactions/<bucket-hash>/<transaction shards>/<transaction-id>/<write-id>`. That deterministic identity is what a provider probe or exact cleanup must bind; it is distinct from the internal transaction-record key.
@@ -638,17 +648,17 @@ The matrix below is the normative approved target, not a blanket description of
| Crash after local transition commit | The exact logical `xl.meta` reference, full recorded source identity (version ID, data directory, modification time, size, and ETag), and remote tuple prove ownership transfer; cleanup only the terminal transaction record |
| Crash after remote DELETE but before journal/free-version cleanup | Retry the same exact idempotent DELETE under the same fences, then conditionally clean local evidence |
| Cancellation | Stop issuing new work, persist monotonic cancellation where the protocol has it, and leave ambiguous durable records for recovery. Cancellation is never rollback proof after authorization |
| Rolling upgrade | Gate writers on the minimum capability required by the format. Known older journal/RPC versions follow their explicit compatibility rule; unknown formats are retained |
| Downgrade | Drain v6 journals and any enabled transition-v2/control protocol before removing their capable workers. Do not write a new format until its downgrade reader behavior and writer gate are specified |
| Rolling upgrade | Gate writers on the minimum capability required by the format or state profile. Transition compaction falls back to the legacy v1 profile until every current topology member answers the exact capability challenge. Known older journal/RPC versions follow their explicit compatibility rule; unknown formats are retained |
| Downgrade | Disable transition compaction admission and drain compact v1 records before removing capable checkpoint validators. Drain v6 journals and any enabled transition-v2/control protocol before removing their capable workers. Do not write a new format until its downgrade reader behavior and writer gate are specified |
| Corrupt or unknown input | Record a diagnosable failure, retain bytes, and block destructive action/completion |
Transition transaction v1, manual job/task/result v1, and receipt v2 do not currently have an implemented persisted-format negotiation for rolling downgrade. The approved transition-v2/control gate above is not current behavior. Until the applicable gate is implemented, caller/operator orchestration must not enable writers whose records required recovery nodes cannot decode. The manual async endpoint does not enforce that fleet gate and a direct request proceeds to job creation. This caller-side fail-closed rule is stricter than treating an unknown record as absent.
The compact transition-transaction v1 state profile has an implemented live homogeneous-fleet gate, but transition v2, manual job/task/result v1, and receipt v2 do not currently have a complete persisted-format negotiation for rolling downgrade. The approved transition-v2/control gate above is not current behavior. Until the applicable gate is implemented, caller/operator orchestration must not enable writers whose records required recovery nodes cannot decode. The manual async endpoint does not enforce that fleet gate and a direct request proceeds to job creation. This caller-side fail-closed rule is stricter than treating an unknown record as absent.
### Current format compatibility decisions
| Family/version | Current reader and writer behavior | Upgrade, downgrade, and ignore rule |
|---|---|---|
| Transition transaction v1 | Writers emit v1; the payload decoder rejects another schema, bad checksum, unknown state, or inconsistent transaction/remote identity. The current record-path parser accepts any shard/extra-component layout and uppercase hex when the final 32-hex UUID parses and matches the payload | There is no intentional ignore path, but exact lowercase canonical-path rejection remains an approved fix. V1 remains the only writer format until the approved v2 fleet gate is implemented; a v2 reader never rewrites an active v1 record |
| Transition transaction v1 | Writers emit v1. A homogeneous live capability permits the compact `UploadOutcomeUnknown@1 -> LocalCommitStarted@2` profile; otherwise writers retain the legacy profile. The payload decoder rejects another schema, bad checksum, unknown state, or inconsistent transaction/remote identity. The current record-path parser accepts any shard/extra-component layout and uppercase hex when the final 32-hex UUID parses and matches the payload | There is no intentional ignore path, but exact lowercase canonical-path rejection remains an approved fix. During rolling upgrade, an unsupported or unavailable peer forces legacy writes. Before downgrade, disable compact admission and drain compact records because older decommission validators may conservatively reject the direct successor. V1 remains the only writer format until the approved v2 fleet gate is implemented; a v2 reader never rewrites an active v1 record |
| Transition transaction v2 and recovery-control/export/disposition v1 | Approved target only; no current reader or writer emits these formats | Roll out read support before the homogeneous writer gate; old readers reject and retain. Disable creation and prove all active records drained before downgrade; never rewrite v2 to v1 |
| Tier mutation intent v1; peer RPC v3/v4 | Durable readers/writers require intent v1. New peers accept signed/canonical v3 and v4 RPC; old v3 peers return an exact authenticated unsupported response to v4 | Pause and drain edit/remove/clear across the mixed interval; do not automatically retry v4 as v3. Unknown durable intent is retained and blocks recovery |
| Manual job/scope/task/result v1 | Writers emit the v1 family. Manual-job runtime recovery accepts an uppercase UUID path when both shard strings match its uppercase prefix, then loads the lowercase canonical job by UUID; the decommission validator recomputes the canonical path and rejects that alias. Other decoder/path/checksum failures stop reconciliation. Runtime capabilities advertise `enqueue_only` and `async`, but the async run handler does not consult a fleet capability gate and a direct request creates a job | Runtime recovery still needs exact lowercase canonical-path validation to prevent alias-driven duplicate work. Caller/operator orchestration must verify every required node and fail closed when capability is unknown or unsupported. An automatic server-side fleet gate and persisted downgrade negotiation remain open; unknown records are never ignored as completed work |
@@ -62,7 +62,6 @@ Object keys are stored as file-system paths under each drive (`{drive}/{bucket}/
| Behavior | RustFS | AWS S3 | Why |
|---|---|---|---|
| Object key with a `.` or `..` path segment, or an empty segment (`//`), such as `a//b/./c/../d` | `400 InvalidArgument` (`check_object_args` in `crates/ecstore/src/bucket/utils.rs`, mirroring MinIO `IsValidObjectPrefix`) | Accepted as an opaque key | A `..` segment would resolve to a parent directory and `.`/`//` segments would alias other keys on disk; encoding them would change the MinIO-compatible on-disk format. |
| Directory marker (key ending in `/`, with or without a body) in a versioned bucket | Stored as the null version: `PutObject`/`HeadObject` report version id `00000000-0000-0000-0000-000000000000`, `ListObjectVersions` reports `null`, and a later PUT of the same key overwrites in place (`put_opts` in `rustfs/src/storage/options.rs`, mirroring MinIO `putOpts`: "for directory objects skip creating new versions") | A real version id per PUT, with a version history | The marker only exists to make an empty prefix listable; keeping a history for it would leave hidden versions behind every prefix delete. Replication still copies the marker as its null version (`test_bucket_replication_replicates_directory_marker_in_versioned_bucket` in `crates/e2e_test/src/replication_extension_test.rs`). |
## Update Rule
+2
View File
@@ -142,6 +142,8 @@ Inspect the aggregate counters before widening scope. Full object-key lists are
Historical transition transactions in `upload_outcome_unknown` state can use an explicit two-stage operator workflow when the tier probe is ambiguous and the provider supports exact version deletion. The endpoint refuses transactions that are still inside their ownership window or are in any other state.
Current fleets can produce two valid v1 state profiles. The legacy profile begins at `upload_started@1` and normally reaches `upload_outcome_unknown@2`, `uploaded`, `local_commit_started`, and `committed`. The compact profile is admitted only while every current member proves `transition_transaction_compaction_v1`; it begins at `upload_outcome_unknown@1` and moves directly to `local_commit_started@2` with a known remote version. Treat `upload_outcome_unknown@1` as a pre-PUT fence, not proof that PUT ran. Treat `local_commit_started@2` as an exact commit fence: if the matching `xl.meta` tuple is complete, recovery removes only the record; otherwise it retains the owner evidence. Do not rewrite either state by hand. An unavailable or older peer automatically makes new transitions use the legacy profile.
1. Inspect the transaction without changing it:
```text
+70 -4
View File
@@ -191,6 +191,7 @@ fn encode_heal_capability_response(
topology_member: &str,
remote_version_state_probe: bool,
recovery_export_probe: bool,
transition_transaction_compaction_probe: bool,
) -> Result<Vec<u8>, Status> {
if recovery_export_probe {
rustfs_protos::encode_remote_version_state_capability(
@@ -198,7 +199,7 @@ fn encode_heal_capability_response(
crate::storage::storage_api::ilm_recovery_export_local_process_epoch().as_bytes(),
)
.map_err(|_| Status::internal("ILM recovery export capability length cannot be represented"))
} else if remote_version_state_probe {
} else if remote_version_state_probe || transition_transaction_compaction_probe {
rustfs_protos::encode_remote_version_state_capability(topology_member, NODE_CAPABILITY_SERVER_EPOCH.as_bytes())
.map_err(|_| Status::internal("remote version state capability length cannot be represented"))
} else {
@@ -1028,7 +1029,13 @@ impl heal_control_service_server::HealControlService for HealControlRpcService {
let remote_version_state_probe = rustfs_protos::is_remote_version_state_capability_probe(&request.get_ref().command);
let cross_pool_fence_probe = rustfs_protos::is_cross_pool_fence_capability_probe(&request.get_ref().command);
let recovery_export_probe = rustfs_protos::is_ilm_recovery_export_capability_probe(&request.get_ref().command);
if remote_version_state_probe || cross_pool_fence_probe || recovery_export_probe {
let transition_transaction_compaction_probe =
rustfs_protos::is_transition_transaction_compaction_capability_probe(&request.get_ref().command);
if remote_version_state_probe
|| cross_pool_fence_probe
|| recovery_export_probe
|| transition_transaction_compaction_probe
{
let topology_member = self
.endpoint_pools()
.await
@@ -1038,7 +1045,12 @@ impl heal_control_service_server::HealControlService for HealControlRpcService {
if topology_member.is_empty() {
return Err(Status::failed_precondition("local topology member identity is unavailable"));
}
let result = encode_heal_capability_response(&topology_member, remote_version_state_probe, recovery_export_probe)?;
let result = encode_heal_capability_response(
&topology_member,
remote_version_state_probe,
recovery_export_probe,
transition_transaction_compaction_probe,
)?;
let canonical_response = rustfs_protos::canonical_heal_control_response_body(
request.get_ref().version,
&request.get_ref().topology_fingerprint,
@@ -4055,9 +4067,50 @@ mod tests {
.expect_err("proof from one challenge must not be reusable");
}
#[tokio::test]
async fn transition_transaction_compaction_probe_authenticates_the_current_topology() {
let _ = rustfs_credentials::set_global_rpc_secret("transition-compaction-node-service-test-secret".to_string());
let endpoints = heal_control_test_endpoints_with_coordinator("node-d", true);
let fingerprint = heal_topology_fingerprint(&endpoints).expect("test topology should hash");
let (service, source) = super::make_heal_control_server_for_source();
*source.write().await = Some(endpoints);
let probe_command = rustfs_protos::transition_transaction_compaction_capability_probe(&[9; 16]);
let mut request = Request::new(HealControlRequest {
version: rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION,
topology_fingerprint: fingerprint.clone(),
command: Bytes::from(probe_command.clone()),
});
let body = rustfs_protos::canonical_heal_control_request_body(
request.get_ref().version,
&request.get_ref().topology_fingerprint,
&request.get_ref().command,
)
.expect("probe should encode");
set_tonic_canonical_body_digest(&mut request, &body).expect("digest metadata should encode");
mark_v2_authenticated(&mut request);
let response = service
.heal_control(request)
.await
.expect("matching topology should admit transition compaction")
.into_inner();
let (member, epoch) = rustfs_protos::decode_remote_version_state_capability(&response.result)
.expect("transition compaction capability should decode");
assert_eq!(member, "node-a:9000");
assert!(!Uuid::from_slice(epoch).expect("capability epoch should be a UUID").is_nil());
let canonical_response = rustfs_protos::canonical_heal_control_response_body(
rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION,
&fingerprint,
&probe_command,
&response.result,
)
.expect("response should encode");
crate::storage::storage_api::verify_tonic_rpc_response_proof(&canonical_response, &response.response_proof)
.expect("outer proof should bind the compaction challenge");
}
#[test]
fn ilm_recovery_export_probe_uses_the_shared_local_process_epoch() {
let result = super::encode_heal_capability_response("node-a:9000", false, true)
let result = super::encode_heal_capability_response("node-a:9000", false, true, false)
.expect("ILM recovery export capability should encode");
let (member, epoch) =
rustfs_protos::decode_remote_version_state_capability(&result).expect("ILM recovery export capability should decode");
@@ -4068,6 +4121,19 @@ mod tests {
);
}
#[test]
fn transition_transaction_compaction_probe_uses_the_node_capability_epoch() {
let remote_version = super::encode_heal_capability_response("node-a:9000", true, false, false)
.expect("remote version capability should encode");
let compaction = super::encode_heal_capability_response("node-a:9000", false, false, true)
.expect("transition transaction compaction capability should encode");
let remote_version = rustfs_protos::decode_remote_version_state_capability(&remote_version)
.expect("remote version capability should decode");
let compaction = rustfs_protos::decode_remote_version_state_capability(&compaction)
.expect("transition transaction compaction capability should decode");
assert_eq!(compaction, remote_version);
}
#[tokio::test]
async fn cross_pool_fence_probe_authenticates_supported_v4_state() {
let _ = rustfs_credentials::set_global_rpc_secret("cross-pool-fence-node-service-test-secret".to_string());