fix(ecstore): make transitioned deletes durable (#5644)

* fix(ecstore): make transitioned deletes durable

* fix(ecstore): journal force deletes

* fix(ecstore): journal force deletes
This commit is contained in:
cxymds
2026-08-03 02:00:11 +08:00
committed by GitHub
parent 8a65017f36
commit 2ce670837c
13 changed files with 823 additions and 264 deletions
+4
View File
@@ -1379,6 +1379,8 @@ mod tests {
backend_identity: Some(identity_a),
version_id_exact: true,
version_state: rustfs_filemeta::TransitionVersionState::Exact,
state: crate::bucket::lifecycle::tier_sweeper::TierDeleteJournalState::Committed,
source: None,
};
let entry_b = Jentry {
obj_name: "remote-b".to_string(),
@@ -1387,6 +1389,8 @@ mod tests {
backend_identity: Some(identity_b),
version_id_exact: true,
version_state: rustfs_filemeta::TransitionVersionState::Exact,
state: crate::bucket::lifecycle::tier_sweeper::TierDeleteJournalState::Committed,
source: None,
};
let remove_a = backend_a.arm_failing_remove_barrier().await;
persist_tier_delete_journal_entry(store_a.clone(), &entry_a)
+316 -3
View File
@@ -13,6 +13,16 @@
// limitations under the License.
use super::*;
use crate::bucket::lifecycle::{
tier_delete_journal::{
abort_prepared_tier_delete_journal_entry as abort_prepared_journal_entry_if_current, commit_tier_delete_journal_entry,
enqueue_committed_tier_delete_journal_entry, persist_tier_delete_journal_entry,
record_tier_delete_journal_backend_identity,
},
tier_sweeper::{
Jentry, attach_tier_delete_source, transitioned_delete_journal_entry_for_source, transitioned_force_delete_journal_entry,
},
};
use crate::bucket::replication::ReplicationObjectBridge;
use crate::disk::OldCurrentSize;
use crate::object_api::DeleteLockFence;
@@ -33,6 +43,181 @@ use std::{
};
use tokio::io::{AsyncRead, ReadBuf};
const FORCE_DELETE_LIST_PAGE_SIZE: i32 = 1_000;
fn build_tier_delete_journal_entry(
bucket: &str,
object: &str,
opts: &ObjectOptions,
source: &ObjectInfo,
) -> Result<Option<Jentry>> {
let version_id = opts.version_id.as_deref().map(Uuid::parse_str).transpose()?;
let source_object = decode_dir_object(object);
let Some(mut je) = (if opts.delete_prefix {
transitioned_force_delete_journal_entry(&source.transitioned_object, source.transition_version_state).map(|mut je| {
attach_tier_delete_source(&mut je, bucket, source_object.as_str(), source, opts.versioned, opts.version_suspended);
je
})
} else {
transitioned_delete_journal_entry_for_source(
version_id,
opts.versioned,
opts.version_suspended,
bucket,
source_object.as_str(),
source,
)
}) else {
return Ok(None);
};
record_tier_delete_journal_backend_identity(&mut je, &source.user_defined).map_err(Error::other)?;
Ok(Some(je))
}
async fn prepare_tier_delete_journal_entry(
api: &Arc<ECStore>,
bucket: &str,
object: &str,
opts: &ObjectOptions,
source: &ObjectInfo,
) -> Result<Option<Jentry>> {
let Some(je) = build_tier_delete_journal_entry(bucket, object, opts, source)? else {
return Ok(None);
};
persist_tier_delete_journal_entry(Arc::clone(api), &je)
.await
.map_err(Error::other)?;
Ok(Some(je))
}
async fn abort_prepared_tier_delete_journal_entry(api: &Arc<ECStore>, je: &Jentry) {
if let Err(err) = abort_prepared_journal_entry_if_current(Arc::clone(api), je).await {
warn!(
object = %je.obj_name,
tier = %je.tier_name,
error = ?err,
"failed to remove aborted tier delete journal"
);
}
}
async fn abort_prepared_tier_delete_journal_entries(api: &Arc<ECStore>, entries: &[Jentry]) {
for entry in entries {
abort_prepared_tier_delete_journal_entry(api, entry).await;
}
}
async fn commit_prepared_tier_delete_journal_entry(api: &Arc<ECStore>, je: &Jentry) {
if let Err(err) = commit_tier_delete_journal_entry(Arc::clone(api), je).await {
warn!(
object = %je.obj_name,
tier = %je.tier_name,
error = ?err,
"tier delete committed locally but journal commit failed; recovery will retry"
);
return;
}
let mut committed = je.clone();
committed.state = crate::bucket::lifecycle::tier_sweeper::TierDeleteJournalState::Committed;
if let Err(err) = enqueue_committed_tier_delete_journal_entry(&committed).await {
warn!(
object = %je.obj_name,
tier = %je.tier_name,
error = ?err,
"tier delete journal committed but could not be queued; recovery will retry"
);
}
}
async fn commit_prepared_tier_delete_journal_entries(api: &Arc<ECStore>, entries: &[Jentry]) {
for entry in entries {
commit_prepared_tier_delete_journal_entry(api, entry).await;
}
}
async fn prepare_prefix_tier_delete_journal_entries(
api: &Arc<ECStore>,
bucket: &str,
prefix: &str,
opts: &ObjectOptions,
) -> Result<Vec<Jentry>> {
let mut marker = None;
let mut version_marker = None;
let mut entries = Vec::new();
loop {
let page = Arc::clone(api)
.list_object_versions_for_lifecycle(
bucket,
prefix,
marker.clone(),
version_marker.clone(),
None,
FORCE_DELETE_LIST_PAGE_SIZE,
)
.await?;
for source in page.objects {
if let Some(entry) = build_tier_delete_journal_entry(bucket, &source.name, opts, &source)? {
entries.push(entry);
}
}
if !page.is_truncated {
break;
}
let next_marker = page
.next_marker
.ok_or_else(|| Error::other("truncated force delete listing has no next marker"))?;
let next_version_marker = page.next_version_idmarker;
if marker.as_deref() == Some(next_marker.as_str()) && version_marker == next_version_marker {
return Err(Error::other("force delete listing marker did not advance"));
}
marker = Some(next_marker);
version_marker = next_version_marker;
}
let mut persisted = Vec::with_capacity(entries.len());
for entry in entries {
if let Err(err) = persist_tier_delete_journal_entry(Arc::clone(api), &entry).await {
abort_prepared_tier_delete_journal_entries(api, &persisted).await;
return Err(Error::other(err));
}
persisted.push(entry);
}
Ok(persisted)
}
async fn delete_prefix_with_tier_delete_journal(
store: &ECStore,
bucket: &str,
object: &str,
opts: &ObjectOptions,
tier_journal_api: Option<&Arc<ECStore>>,
) -> Result<()> {
let journal_entry = if let Some(api) = tier_journal_api {
Some(prepare_prefix_tier_delete_journal_entries(api, bucket, object, opts).await?)
} else {
None
};
let result = store.delete_prefix(bucket, object, opts).await;
match result {
Ok(()) => {
if let (Some(api), Some(entries)) = (tier_journal_api, journal_entry.as_ref()) {
commit_prepared_tier_delete_journal_entries(api, entries).await;
}
Ok(())
}
Err(err) => {
if let (Some(api), Some(entries)) = (tier_journal_api, journal_entry.as_ref()) {
abort_prepared_tier_delete_journal_entries(api, entries).await;
}
Err(err)
}
}
}
/// A GET whose object identity has been resolved while its namespace read lock
/// remains held, but whose body reader has not been constructed yet.
///
@@ -1092,8 +1277,49 @@ impl ECStore {
purged
}
pub async fn delete_object_with_tier_delete_journal(
self: &Arc<Self>,
bucket: &str,
object: &str,
opts: ObjectOptions,
) -> Result<ObjectInfo> {
let result = self
.handle_delete_object_with_journal(bucket, object, opts, Some(Arc::clone(self)))
.await;
if result.is_ok() {
list_objects::observe_list_objects_mutation(self, bucket).await;
}
result
}
pub async fn delete_objects_with_tier_delete_journal(
self: &Arc<Self>,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
let result = self
.handle_delete_objects_with_journal(bucket, objects, opts, Some(Arc::clone(self)))
.await;
let success_count = result.1.iter().filter(|err| err.is_none()).count();
if success_count > 0 {
list_objects::observe_list_objects_mutations(self, bucket, success_count).await;
}
result
}
#[instrument(skip(self))]
pub(super) async fn handle_delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result<ObjectInfo> {
self.handle_delete_object_with_journal(bucket, object, opts, None).await
}
pub(super) async fn handle_delete_object_with_journal(
&self,
bucket: &str,
object: &str,
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> Result<ObjectInfo> {
check_del_obj_args(bucket, object)?;
let object = if opts.delete_prefix && !opts.delete_prefix_object {
@@ -1103,11 +1329,12 @@ impl ECStore {
};
let object = object.as_str();
let mut opts = opts;
opts.tier_delete_journal_api = tier_journal_api.clone();
if opts.delete_prefix && !opts.delete_prefix_object {
// Prefix deletes cover multiple object keys; an exact lock on the prefix string
// would not protect child objects.
self.delete_prefix(bucket, object, &opts).await?;
delete_prefix_with_tier_delete_journal(self, bucket, object, &opts, tier_journal_api.as_ref()).await?;
return Ok(ObjectInfo::default());
}
@@ -1119,7 +1346,7 @@ impl ECStore {
};
if opts.delete_prefix {
self.delete_prefix(bucket, object, &opts).await?;
delete_prefix_with_tier_delete_journal(self, bucket, object, &opts, tier_journal_api.as_ref()).await?;
return Ok(ObjectInfo::default());
}
@@ -1215,8 +1442,25 @@ impl ECStore {
));
}
let journal_entry = if let Some(api) = tier_journal_api.as_ref() {
prepare_tier_delete_journal_entry(api, bucket, object, &opts, &pinfo.object_info).await?
} else {
None
};
if !errs.is_empty() && !opts.versioned && !opts.version_suspended {
let mut obj = self.delete_object_from_all_pools(bucket, object, &opts, errs).await?;
let mut obj = match self.delete_object_from_all_pools(bucket, object, &opts, errs).await {
Ok(obj) => obj,
Err(err) => {
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
abort_prepared_tier_delete_journal_entry(api, je).await;
}
return Err(err);
}
};
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
commit_prepared_tier_delete_journal_entry(api, je).await;
}
obj.name = decode_dir_object(object);
return Ok(obj);
}
@@ -1224,6 +1468,9 @@ impl ECStore {
for pool in self.pools.iter() {
match pool.delete_object(bucket, object, opts.clone()).await {
Ok(res) => {
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
commit_prepared_tier_delete_journal_entry(api, je).await;
}
let mut obj = res;
obj.name = decode_dir_object(object);
return Ok(obj);
@@ -1236,6 +1483,10 @@ impl ECStore {
}
}
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
abort_prepared_tier_delete_journal_entry(api, je).await;
}
if let Some(ver) = opts.version_id {
return Err(StorageError::VersionNotFound(bucket.to_owned(), object.to_owned(), ver));
}
@@ -1249,6 +1500,16 @@ impl ECStore {
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
self.handle_delete_objects_with_journal(bucket, objects, opts, None).await
}
pub(super) async fn handle_delete_objects_with_journal(
&self,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
// encode object name
let objects: Vec<ObjectToDelete> = objects
@@ -1269,6 +1530,7 @@ impl ECStore {
}
let mut opts = opts;
opts.tier_delete_journal_api = tier_journal_api;
if opts.delete_replication_config_snapshot.is_none() {
match ReplicationObjectBridge::delete_request_config_in(&self.ctx, bucket).await {
Ok(snapshot) => opts.delete_replication_config_snapshot = Some(Arc::new(snapshot)),
@@ -1625,6 +1887,7 @@ impl ECStore {
mod tests {
use super::*;
use crate::bucket::lifecycle::core::TRANSITION_COMPLETE;
use crate::bucket::lifecycle::tier_sweeper::TierDeleteJournalState;
use crate::bucket::replication::{
ReplicationState, ReplicationStatusType, VersionPurgeStatusType, replication_state_to_filemeta, replication_statuses_map,
version_purge_statuses_map,
@@ -1640,6 +1903,7 @@ mod tests {
};
use crate::set_disk::SetDisks;
use crate::storage_api_contracts::bucket::MakeBucketOptions;
use crate::storage_api_contracts::lifecycle::TransitionedObject;
use bytes::Bytes;
use std::io::Cursor;
use std::sync::Arc;
@@ -1660,6 +1924,55 @@ mod tests {
struct BodyCacheHookGuard;
#[test]
fn tier_delete_entry_is_prepared_and_bound_to_source_generation() {
let identity = [9_u8; 32];
let mut metadata = HashMap::new();
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex(identity),
);
let version_id = Uuid::from_u128(1);
let data_dir = Uuid::from_u128(2);
let source = ObjectInfo {
bucket: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(version_id),
data_dir: Some(data_dir),
user_defined: Arc::new(metadata),
transitioned_object: TransitionedObject {
name: "remote/object".to_string(),
version_id: "remote-version".to_string(),
tier: "WARM".to_string(),
status: TRANSITION_COMPLETE.to_string(),
..Default::default()
},
transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
..Default::default()
};
let entry = build_tier_delete_journal_entry(
"bucket",
"object",
&ObjectOptions {
version_id: Some(version_id.to_string()),
versioned: true,
..Default::default()
},
&source,
)
.expect("transition source should produce a journal entry")
.expect("completed transition should be journaled");
assert_eq!(entry.state, TierDeleteJournalState::Prepared);
assert_eq!(entry.backend_identity, Some(identity));
let data_dir_string = data_dir.to_string();
assert_eq!(
entry.source.as_ref().and_then(|source| source.data_dir.as_deref()),
Some(data_dir_string.as_str())
);
}
impl Drop for BodyCacheHookGuard {
fn drop(&mut self) {
clear_get_object_body_cache_hook();