mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 04:16:38 +00:00
Compare commits
15 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2f048466b4 | |||
| 2a4021c3c5 | |||
| c6ea2dd9b3 | |||
| 83b87d2a3e | |||
| a61a95ae48 | |||
| fd3b47f92b | |||
| fe1d5c7b5d | |||
| 9fa3e9c41c | |||
| 1fec344f43 | |||
| acffa7511f | |||
| 9a0cd48f9e | |||
| 15b8e8860f | |||
| 812cf9662a | |||
| 50c208f716 | |||
| 44351323e5 |
@@ -78,6 +78,12 @@ test-group = 'embedded-test-ports'
|
|||||||
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
|
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
|
||||||
test-group = 'ecstore-serial-flaky'
|
test-group = 'ecstore-serial-flaky'
|
||||||
|
|
||||||
|
# The durable ILM decommission regressions build isolated multi-pool stores and
|
||||||
|
# deliberately take source or target disks offline while checking fencing.
|
||||||
|
[[profile.default.overrides]]
|
||||||
|
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
|
||||||
|
test-group = 'ecstore-serial-flaky'
|
||||||
|
|
||||||
# Serialize the bucket-incarnation / lifecycle-fence tests. They drive
|
# Serialize the bucket-incarnation / lifecycle-fence tests. They drive
|
||||||
# init_bucket_metadata_sys and bucket_metadata_sys_of, i.e. process-global
|
# init_bucket_metadata_sys and bucket_metadata_sys_of, i.e. process-global
|
||||||
# OnceLock state that serial_test's #[serial] cannot protect across nextest's
|
# OnceLock state that serial_test's #[serial] cannot protect across nextest's
|
||||||
@@ -190,6 +196,10 @@ test-group = 'embedded-test-ports'
|
|||||||
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
|
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
|
||||||
test-group = 'ecstore-serial-flaky'
|
test-group = 'ecstore-serial-flaky'
|
||||||
|
|
||||||
|
[[profile.ci.overrides]]
|
||||||
|
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
|
||||||
|
test-group = 'ecstore-serial-flaky'
|
||||||
|
|
||||||
# Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile
|
# Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile
|
||||||
# too (see the matching default-profile override near the top). No retries.
|
# too (see the matching default-profile override near the top). No retries.
|
||||||
[[profile.ci.overrides]]
|
[[profile.ci.overrides]]
|
||||||
|
|||||||
@@ -9136,6 +9136,7 @@ mod tests {
|
|||||||
assert_eq!(loaded.report.scanned, 37);
|
assert_eq!(loaded.report.scanned, 37);
|
||||||
assert_eq!(loaded.report.eligible, 11);
|
assert_eq!(loaded.report.eligible, 11);
|
||||||
assert_eq!(loaded.report.enqueued, 5);
|
assert_eq!(loaded.report.enqueued, 5);
|
||||||
|
assert_eq!(loaded.cursor_revision, Some(1));
|
||||||
assert!(loaded.lease_expires_at_unix_nanos > 0);
|
assert!(loaded.lease_expires_at_unix_nanos > 0);
|
||||||
let token = loaded
|
let token = loaded
|
||||||
.report
|
.report
|
||||||
@@ -9154,6 +9155,30 @@ mod tests {
|
|||||||
assert_eq!(admission.lease_id, loaded.lease_id);
|
assert_eq!(admission.lease_id, loaded.lease_id);
|
||||||
assert_eq!(admission.lease_expires_at_unix_nanos, loaded.lease_expires_at_unix_nanos);
|
assert_eq!(admission.lease_expires_at_unix_nanos, loaded.lease_expires_at_unix_nanos);
|
||||||
|
|
||||||
|
let mut same_marker_report = report.clone();
|
||||||
|
same_marker_report.scanned += 1;
|
||||||
|
persist_manual_transition_page_checkpoint(
|
||||||
|
&checkpoint_options,
|
||||||
|
&same_marker_report,
|
||||||
|
Some("logs/page-end".to_string()),
|
||||||
|
Some("opaque-next-version".to_string()),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("same-marker version checkpoint should persist through the durable progress sink");
|
||||||
|
let same_marker_checkpointed = load_manual_transition_job_record(ecstore.clone(), job_id)
|
||||||
|
.await
|
||||||
|
.expect("same-marker version checkpoint should reload");
|
||||||
|
assert_eq!(same_marker_checkpointed.cursor_revision, Some(2));
|
||||||
|
let (_, version_marker) = decode_manual_transition_continuation_token(
|
||||||
|
same_marker_checkpointed
|
||||||
|
.report
|
||||||
|
.continuation_token
|
||||||
|
.as_deref()
|
||||||
|
.expect("same-marker version checkpoint should persist a cursor"),
|
||||||
|
)
|
||||||
|
.expect("same-marker version cursor should decode");
|
||||||
|
assert_eq!(version_marker.as_deref(), Some("opaque-next-version"));
|
||||||
|
|
||||||
create_test_bucket(&ecstore, &bucket).await;
|
create_test_bucket(&ecstore, &bucket).await;
|
||||||
let lifecycle_xml = format!(
|
let lifecycle_xml = format!(
|
||||||
r#"<?xml version="1.0" encoding="UTF-8"?>
|
r#"<?xml version="1.0" encoding="UTF-8"?>
|
||||||
@@ -9210,6 +9235,7 @@ mod tests {
|
|||||||
assert_eq!(checkpointed.report.scanned, 1000);
|
assert_eq!(checkpointed.report.scanned, 1000);
|
||||||
assert_eq!(checkpointed.report.eligible, 1000);
|
assert_eq!(checkpointed.report.eligible, 1000);
|
||||||
assert_eq!(checkpointed.report.dry_run_eligible, 1000);
|
assert_eq!(checkpointed.report.dry_run_eligible, 1000);
|
||||||
|
assert_eq!(checkpointed.cursor_revision, Some(3));
|
||||||
let token = checkpointed
|
let token = checkpointed
|
||||||
.report
|
.report
|
||||||
.continuation_token
|
.continuation_token
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
@@ -24,6 +24,10 @@ use crate::bucket::lifecycle::bucket_lifecycle_ops::{
|
|||||||
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport,
|
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport,
|
||||||
};
|
};
|
||||||
use crate::bucket::lifecycle::config_boundary;
|
use crate::bucket::lifecycle::config_boundary;
|
||||||
|
use crate::bucket::lifecycle::durable_namespace::{
|
||||||
|
MANUAL_TRANSITION_JOB_NAMESPACE, MANUAL_TRANSITION_SCOPE_NAMESPACE, MANUAL_TRANSITION_TASK_NAMESPACE,
|
||||||
|
MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE,
|
||||||
|
};
|
||||||
use crate::disk::RUSTFS_META_BUCKET;
|
use crate::disk::RUSTFS_META_BUCKET;
|
||||||
use crate::error::{Error, Result as EcstoreResult};
|
use crate::error::{Error, Result as EcstoreResult};
|
||||||
use crate::object_api::ObjectOptions;
|
use crate::object_api::ObjectOptions;
|
||||||
@@ -34,10 +38,10 @@ use crate::store::ECStore;
|
|||||||
pub const MANUAL_TRANSITION_JOB_SCHEMA: &str = "rustfs-manual-transition-job-v1";
|
pub const MANUAL_TRANSITION_JOB_SCHEMA: &str = "rustfs-manual-transition-job-v1";
|
||||||
pub const MANUAL_TRANSITION_TASK_SCHEMA: &str = "rustfs-manual-transition-task-v1";
|
pub const MANUAL_TRANSITION_TASK_SCHEMA: &str = "rustfs-manual-transition-task-v1";
|
||||||
pub const MANUAL_TRANSITION_WORKER_RESULT_SCHEMA: &str = "rustfs-manual-transition-worker-result-v1";
|
pub const MANUAL_TRANSITION_WORKER_RESULT_SCHEMA: &str = "rustfs-manual-transition-worker-result-v1";
|
||||||
pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = "ilm/manual-transition/jobs";
|
pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = MANUAL_TRANSITION_JOB_NAMESPACE.prefix;
|
||||||
pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = "ilm/manual-transition/scopes";
|
pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = MANUAL_TRANSITION_SCOPE_NAMESPACE.prefix;
|
||||||
pub const MANUAL_TRANSITION_TASK_PREFIX: &str = "ilm/manual-transition/tasks";
|
pub const MANUAL_TRANSITION_TASK_PREFIX: &str = MANUAL_TRANSITION_TASK_NAMESPACE.prefix;
|
||||||
pub const MANUAL_TRANSITION_WORKER_RESULT_PREFIX: &str = "ilm/manual-transition/results";
|
pub const MANUAL_TRANSITION_WORKER_RESULT_PREFIX: &str = MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE.prefix;
|
||||||
pub const MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE: usize = 64 * 1024;
|
pub const MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE: usize = 64 * 1024;
|
||||||
pub const MAX_MANUAL_TRANSITION_TASK_RECORD_SIZE: usize = 16 * 1024;
|
pub const MAX_MANUAL_TRANSITION_TASK_RECORD_SIZE: usize = 16 * 1024;
|
||||||
pub const MAX_MANUAL_TRANSITION_WORKER_RESULT_RECORD_SIZE: usize = 8 * 1024;
|
pub const MAX_MANUAL_TRANSITION_WORKER_RESULT_RECORD_SIZE: usize = 8 * 1024;
|
||||||
@@ -195,6 +199,8 @@ pub struct ManualTransitionJobRecord {
|
|||||||
pub updated_at_unix_nanos: i128,
|
pub updated_at_unix_nanos: i128,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub completed_at_unix_nanos: Option<i128>,
|
pub completed_at_unix_nanos: Option<i128>,
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
pub cursor_revision: Option<u64>,
|
||||||
pub report: ManualTransitionRunReport,
|
pub report: ManualTransitionRunReport,
|
||||||
pub queue_snapshot: ManualTransitionQueueSnapshot,
|
pub queue_snapshot: ManualTransitionQueueSnapshot,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
@@ -224,6 +230,7 @@ impl ManualTransitionJobRecord {
|
|||||||
created_at_unix_nanos: now,
|
created_at_unix_nanos: now,
|
||||||
updated_at_unix_nanos: now,
|
updated_at_unix_nanos: now,
|
||||||
completed_at_unix_nanos: None,
|
completed_at_unix_nanos: None,
|
||||||
|
cursor_revision: Some(0),
|
||||||
report: ManualTransitionRunReport {
|
report: ManualTransitionRunReport {
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
prefix: options.prefix.clone(),
|
prefix: options.prefix.clone(),
|
||||||
@@ -238,7 +245,7 @@ impl ManualTransitionJobRecord {
|
|||||||
|
|
||||||
pub fn complete(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
|
pub fn complete(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||||
self.scan_completed = true;
|
self.scan_completed = true;
|
||||||
self.report.merge_scan_report_preserving_worker(&report);
|
self.merge_scan_report(&report);
|
||||||
self.queue_snapshot = queue_snapshot;
|
self.queue_snapshot = queue_snapshot;
|
||||||
self.error = None;
|
self.error = None;
|
||||||
self.mark_terminal_if_worker_drained();
|
self.mark_terminal_if_worker_drained();
|
||||||
@@ -316,7 +323,7 @@ impl ManualTransitionJobRecord {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
self.queue_snapshot = queue_snapshot;
|
self.queue_snapshot = queue_snapshot;
|
||||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
self.advance_updated_at();
|
||||||
self.mark_terminal_if_worker_drained();
|
self.mark_terminal_if_worker_drained();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -359,14 +366,14 @@ impl ManualTransitionJobRecord {
|
|||||||
self.report.tier_failure = scan_tier_failure.saturating_add(transition_failed);
|
self.report.tier_failure = scan_tier_failure.saturating_add(transition_failed);
|
||||||
self.report.tier_failure_by_reason = scan_tier_failure_by_reason;
|
self.report.tier_failure_by_reason = scan_tier_failure_by_reason;
|
||||||
self.queue_snapshot = queue_snapshot;
|
self.queue_snapshot = queue_snapshot;
|
||||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
self.advance_updated_at();
|
||||||
self.mark_terminal_if_worker_drained();
|
self.mark_terminal_if_worker_drained();
|
||||||
true
|
true
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn mark_cancel_requested(&mut self) {
|
pub fn mark_cancel_requested(&mut self) {
|
||||||
self.cancel_requested = true;
|
self.cancel_requested = true;
|
||||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
self.advance_updated_at();
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn claim_recovery_lease(&mut self, owner_id: impl Into<String>, queue_snapshot: ManualTransitionQueueSnapshot) {
|
pub fn claim_recovery_lease(&mut self, owner_id: impl Into<String>, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||||
@@ -380,7 +387,7 @@ impl ManualTransitionJobRecord {
|
|||||||
pub fn abandon_recovery_lease(&mut self, lease_id: Uuid) {
|
pub fn abandon_recovery_lease(&mut self, lease_id: Uuid) {
|
||||||
if self.state == ManualTransitionJobState::Running && self.lease_id == lease_id {
|
if self.state == ManualTransitionJobState::Running && self.lease_id == lease_id {
|
||||||
self.lease_expires_at_unix_nanos = 0;
|
self.lease_expires_at_unix_nanos = 0;
|
||||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
self.advance_updated_at();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -398,7 +405,7 @@ impl ManualTransitionJobRecord {
|
|||||||
|
|
||||||
pub fn renew_lease(&mut self, queue_snapshot: ManualTransitionQueueSnapshot) {
|
pub fn renew_lease(&mut self, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||||
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||||
self.updated_at_unix_nanos = now;
|
self.updated_at_unix_nanos = self.updated_at_unix_nanos.saturating_add(1).max(now);
|
||||||
self.lease_expires_at_unix_nanos = manual_transition_job_lease_expires_at(now);
|
self.lease_expires_at_unix_nanos = manual_transition_job_lease_expires_at(now);
|
||||||
self.queue_snapshot = queue_snapshot;
|
self.queue_snapshot = queue_snapshot;
|
||||||
}
|
}
|
||||||
@@ -438,11 +445,18 @@ impl ManualTransitionJobRecord {
|
|||||||
|
|
||||||
pub fn update_running_progress(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
|
pub fn update_running_progress(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||||
if self.state == ManualTransitionJobState::Running {
|
if self.state == ManualTransitionJobState::Running {
|
||||||
self.report.merge_scan_report_preserving_worker(&report);
|
self.merge_scan_report(&report);
|
||||||
self.renew_lease(queue_snapshot);
|
self.renew_lease(queue_snapshot);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn merge_scan_report(&mut self, report: &ManualTransitionRunReport) {
|
||||||
|
if self.report.continuation_token != report.continuation_token {
|
||||||
|
self.cursor_revision = Some(self.cursor_revision.unwrap_or(0).saturating_add(1));
|
||||||
|
}
|
||||||
|
self.report.merge_scan_report_preserving_worker(report);
|
||||||
|
}
|
||||||
|
|
||||||
pub fn mark_unknown_if_unowned(&mut self) {
|
pub fn mark_unknown_if_unowned(&mut self) {
|
||||||
if self.state == ManualTransitionJobState::Running {
|
if self.state == ManualTransitionJobState::Running {
|
||||||
self.state = ManualTransitionJobState::Unknown;
|
self.state = ManualTransitionJobState::Unknown;
|
||||||
@@ -463,9 +477,13 @@ impl ManualTransitionJobRecord {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn mark_updated_terminal(&mut self) {
|
fn mark_updated_terminal(&mut self) {
|
||||||
|
self.advance_updated_at();
|
||||||
|
self.completed_at_unix_nanos = Some(self.updated_at_unix_nanos);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn advance_updated_at(&mut self) {
|
||||||
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||||
self.updated_at_unix_nanos = now;
|
self.updated_at_unix_nanos = self.updated_at_unix_nanos.saturating_add(1).max(now);
|
||||||
self.completed_at_unix_nanos = Some(now);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn mark_terminal_if_worker_drained(&mut self) {
|
fn mark_terminal_if_worker_drained(&mut self) {
|
||||||
@@ -1109,7 +1127,8 @@ pub fn manual_transition_scope_record_object_name(scope_key: &str) -> Result<Str
|
|||||||
pub async fn save_manual_transition_job_record(api: Arc<ECStore>, job: &ManualTransitionJobRecord) -> EcstoreResult<()> {
|
pub async fn save_manual_transition_job_record(api: Arc<ECStore>, job: &ManualTransitionJobRecord) -> EcstoreResult<()> {
|
||||||
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
|
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
|
||||||
let data = job.encode().map_err(manual_transition_job_store_error)?;
|
let data = job.encode().map_err(manual_transition_job_store_error)?;
|
||||||
config_boundary::save_config(api, &object, data).await
|
config_boundary::save_config(api.clone(), &object, data.clone()).await?;
|
||||||
|
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn load_manual_transition_job_record(api: Arc<ECStore>, job_id: Uuid) -> EcstoreResult<ManualTransitionJobRecord> {
|
pub async fn load_manual_transition_job_record(api: Arc<ECStore>, job_id: Uuid) -> EcstoreResult<ManualTransitionJobRecord> {
|
||||||
@@ -1142,9 +1161,9 @@ pub async fn save_manual_transition_job_record_if_current(
|
|||||||
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
|
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
|
||||||
let data = job.encode().map_err(manual_transition_job_store_error)?;
|
let data = job.encode().map_err(manual_transition_job_store_error)?;
|
||||||
config_boundary::save_config_with_opts_quiet(
|
config_boundary::save_config_with_opts_quiet(
|
||||||
api,
|
api.clone(),
|
||||||
&object,
|
&object,
|
||||||
data,
|
data.clone(),
|
||||||
&ObjectOptions {
|
&ObjectOptions {
|
||||||
max_parity: true,
|
max_parity: true,
|
||||||
http_preconditions: Some(HTTPPreconditions {
|
http_preconditions: Some(HTTPPreconditions {
|
||||||
@@ -1154,7 +1173,8 @@ pub async fn save_manual_transition_job_record_if_current(
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await?;
|
||||||
|
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Applies a job-record mutation with optimistic concurrency control.
|
/// Applies a job-record mutation with optimistic concurrency control.
|
||||||
@@ -1592,9 +1612,9 @@ pub async fn save_manual_transition_scope_admission_if_absent(
|
|||||||
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
|
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
|
||||||
let data = serde_json::to_vec(admission).map_err(Error::other)?;
|
let data = serde_json::to_vec(admission).map_err(Error::other)?;
|
||||||
config_boundary::save_config_with_opts(
|
config_boundary::save_config_with_opts(
|
||||||
api,
|
api.clone(),
|
||||||
&object,
|
&object,
|
||||||
data,
|
data.clone(),
|
||||||
&ObjectOptions {
|
&ObjectOptions {
|
||||||
max_parity: true,
|
max_parity: true,
|
||||||
http_preconditions: Some(HTTPPreconditions {
|
http_preconditions: Some(HTTPPreconditions {
|
||||||
@@ -1604,7 +1624,8 @@ pub async fn save_manual_transition_scope_admission_if_absent(
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await?;
|
||||||
|
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn load_manual_transition_scope_admission(
|
pub async fn load_manual_transition_scope_admission(
|
||||||
@@ -1642,9 +1663,9 @@ pub async fn save_manual_transition_scope_admission_if_current(
|
|||||||
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
|
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
|
||||||
let data = serde_json::to_vec(admission).map_err(Error::other)?;
|
let data = serde_json::to_vec(admission).map_err(Error::other)?;
|
||||||
match config_boundary::save_config_with_opts(
|
match config_boundary::save_config_with_opts(
|
||||||
api,
|
api.clone(),
|
||||||
&object,
|
&object,
|
||||||
data,
|
data.clone(),
|
||||||
&ObjectOptions {
|
&ObjectOptions {
|
||||||
max_parity: true,
|
max_parity: true,
|
||||||
http_preconditions: Some(HTTPPreconditions {
|
http_preconditions: Some(HTTPPreconditions {
|
||||||
@@ -1660,7 +1681,8 @@ pub async fn save_manual_transition_scope_admission_if_current(
|
|||||||
Err(Error::PreconditionFailed)
|
Err(Error::PreconditionFailed)
|
||||||
}
|
}
|
||||||
result => result,
|
result => result,
|
||||||
}
|
}?;
|
||||||
|
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn claim_manual_transition_scope_admission(
|
pub async fn claim_manual_transition_scope_admission(
|
||||||
@@ -1953,13 +1975,15 @@ pub async fn delete_manual_transition_scope_admission_if_current(
|
|||||||
job_id: Uuid,
|
job_id: Uuid,
|
||||||
lease_id: Uuid,
|
lease_id: Uuid,
|
||||||
) -> EcstoreResult<bool> {
|
) -> EcstoreResult<bool> {
|
||||||
let etag = match load_manual_transition_scope_admission_with_etag(api.clone(), scope_key).await {
|
let (admission, etag) = match load_manual_transition_scope_admission_with_etag(api.clone(), scope_key).await {
|
||||||
Ok((admission, etag)) if admission.job_id == job_id && admission.lease_id == lease_id => etag,
|
Ok((admission, etag)) if admission.job_id == job_id && admission.lease_id == lease_id => (admission, etag),
|
||||||
Ok(_) => return Ok(false),
|
Ok(_) => return Ok(false),
|
||||||
Err(Error::ConfigNotFound) => return Ok(true),
|
Err(Error::ConfigNotFound) => return Ok(true),
|
||||||
Err(err) => return Err(err),
|
Err(err) => return Err(err),
|
||||||
};
|
};
|
||||||
let object = manual_transition_scope_record_object_name(scope_key).map_err(manual_transition_job_store_error)?;
|
let object = manual_transition_scope_record_object_name(scope_key).map_err(manual_transition_job_store_error)?;
|
||||||
|
let data = serde_json::to_vec(&admission).map_err(Error::other)?;
|
||||||
|
api.record_durable_ilm_decommission_terminal(&object, &data).await?;
|
||||||
match config_boundary::delete_config_if_match(api, &object, &etag).await {
|
match config_boundary::delete_config_if_match(api, &object, &etag).await {
|
||||||
Ok(()) | Err(Error::ConfigNotFound) => Ok(true),
|
Ok(()) | Err(Error::ConfigNotFound) => Ok(true),
|
||||||
Err(Error::PreconditionFailed) => Ok(false),
|
Err(Error::PreconditionFailed) => Ok(false),
|
||||||
@@ -2577,6 +2601,25 @@ mod tests {
|
|||||||
assert!(decoded.report.tier_failure_by_reason.is_empty());
|
assert!(decoded.report.tier_failure_by_reason.is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn manual_transition_job_record_decodes_legacy_cursor_without_revision() {
|
||||||
|
let options = ManualTransitionRunOptions::default();
|
||||||
|
let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
|
||||||
|
let encoded = record.encode().expect("job record should encode");
|
||||||
|
let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded job should be json");
|
||||||
|
value["job"]
|
||||||
|
.as_object_mut()
|
||||||
|
.expect("job should be object")
|
||||||
|
.remove("cursor_revision");
|
||||||
|
let record_bytes = serde_json::to_vec(&value["job"]).expect("legacy job should encode");
|
||||||
|
value["content_sha256"] = serde_json::Value::String(hex_sha256(&record_bytes, ToOwned::to_owned));
|
||||||
|
let legacy = serde_json::to_vec(&value).expect("legacy envelope should encode");
|
||||||
|
|
||||||
|
let decoded = ManualTransitionJobRecord::decode(record.job_id, &legacy).expect("legacy job should decode");
|
||||||
|
|
||||||
|
assert_eq!(decoded.cursor_revision, None);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn manual_transition_job_record_rejects_unknown_report_fields() {
|
fn manual_transition_job_record_rejects_unknown_report_fields() {
|
||||||
let options = ManualTransitionRunOptions::default();
|
let options = ManualTransitionRunOptions::default();
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ pub mod bucket_lifecycle_audit;
|
|||||||
pub mod bucket_lifecycle_ops;
|
pub mod bucket_lifecycle_ops;
|
||||||
mod config_boundary;
|
mod config_boundary;
|
||||||
pub mod core;
|
pub mod core;
|
||||||
|
mod durable_namespace;
|
||||||
pub mod evaluator;
|
pub mod evaluator;
|
||||||
pub mod manual_transition_job;
|
pub mod manual_transition_job;
|
||||||
mod metadata_boundary;
|
mod metadata_boundary;
|
||||||
@@ -31,3 +32,8 @@ pub mod tier_free_version_recovery;
|
|||||||
pub mod tier_last_day_stats;
|
pub mod tier_last_day_stats;
|
||||||
pub mod tier_sweeper;
|
pub mod tier_sweeper;
|
||||||
pub mod transition_transaction;
|
pub mod transition_transaction;
|
||||||
|
|
||||||
|
pub(crate) use durable_namespace::{
|
||||||
|
DurableIlmRecordCheckpoint, ILM_META_PREFIX, ValidatedDurableIlmRecord, classify_durable_ilm_record,
|
||||||
|
validate_durable_ilm_record,
|
||||||
|
};
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ use tokio_util::sync::CancellationToken;
|
|||||||
use tracing::{debug, warn};
|
use tracing::{debug, warn};
|
||||||
|
|
||||||
use crate::bucket::lifecycle::config_boundary;
|
use crate::bucket::lifecycle::config_boundary;
|
||||||
|
use crate::bucket::lifecycle::durable_namespace::TIER_DELETE_JOURNAL_NAMESPACE;
|
||||||
use crate::bucket::lifecycle::runtime_boundary;
|
use crate::bucket::lifecycle::runtime_boundary;
|
||||||
use crate::bucket::lifecycle::tier_sweeper::{
|
use crate::bucket::lifecycle::tier_sweeper::{
|
||||||
Jentry, TierDeleteJournalState, TierDeleteSourceIdentity,
|
Jentry, TierDeleteJournalState, TierDeleteSourceIdentity,
|
||||||
@@ -49,7 +50,7 @@ const TIER_DELETE_JOURNAL_VERSION: u8 = 2;
|
|||||||
const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3;
|
const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3;
|
||||||
const TIER_DELETE_JOURNAL_STATE_VERSION: u8 = 4;
|
const TIER_DELETE_JOURNAL_STATE_VERSION: u8 = 4;
|
||||||
const TIER_DELETE_JOURNAL_TRANSACTION_VERSION: u8 = 5;
|
const TIER_DELETE_JOURNAL_TRANSACTION_VERSION: u8 = 5;
|
||||||
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/";
|
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = TIER_DELETE_JOURNAL_NAMESPACE.prefix;
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
#[serde(deny_unknown_fields)]
|
#[serde(deny_unknown_fields)]
|
||||||
@@ -432,6 +433,21 @@ async fn process_committed_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jen
|
|||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
|
let path = tier_delete_journal_object_name(je);
|
||||||
|
let data = encode_tier_delete_journal_entry(je).map_err(std::io::Error::other)?;
|
||||||
|
let target_pool_indices = api
|
||||||
|
.record_durable_ilm_decommission_terminal_target_pools(&path, &data)
|
||||||
|
.await
|
||||||
|
.map_err(std::io::Error::other)?;
|
||||||
|
if let Some(target_pool_indices) = target_pool_indices {
|
||||||
|
for target_pool_idx in target_pool_indices {
|
||||||
|
match config_boundary::delete_config(api.pools[target_pool_idx].clone(), &path).await {
|
||||||
|
Ok(()) | Err(Error::ConfigNotFound) => {}
|
||||||
|
Err(err) => return Err(std::io::Error::other(err)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
remove_tier_delete_journal_entry(api, je).await
|
remove_tier_delete_journal_entry(api, je).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ use tracing::{debug, warn};
|
|||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
use crate::bucket::lifecycle::config_boundary;
|
use crate::bucket::lifecycle::config_boundary;
|
||||||
|
use crate::bucket::lifecycle::durable_namespace::TRANSITION_TRANSACTION_NAMESPACE;
|
||||||
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
||||||
use crate::bucket::lifecycle::tier_sweeper::{
|
use crate::bucket::lifecycle::tier_sweeper::{
|
||||||
delete_confirmed_transition_candidate_exact_with_lease_idempotent,
|
delete_confirmed_transition_candidate_exact_with_lease_idempotent,
|
||||||
@@ -42,7 +43,7 @@ const TRANSITION_TRANSACTION_RECOVERY_INTERVAL: Duration = Duration::from_secs(6
|
|||||||
const TRANSITION_TRANSACTION_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
|
const TRANSITION_TRANSACTION_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
|
||||||
pub const TRANSITION_TRANSACTION_SCHEMA: &str = "rustfs-transition-transaction-v1";
|
pub const TRANSITION_TRANSACTION_SCHEMA: &str = "rustfs-transition-transaction-v1";
|
||||||
pub const TRANSITION_TRANSACTION_PREFIX: &str = "ilm/transition-transactions";
|
pub const TRANSITION_TRANSACTION_PREFIX: &str = "ilm/transition-transactions";
|
||||||
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = "ilm/transition-transactions/records";
|
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = TRANSITION_TRANSACTION_NAMESPACE.prefix;
|
||||||
pub const MAX_TRANSITION_TRANSACTION_SIZE: usize = 64 * 1024;
|
pub const MAX_TRANSITION_TRANSACTION_SIZE: usize = 64 * 1024;
|
||||||
|
|
||||||
pub type Result<T> = std::result::Result<T, TransitionTransactionError>;
|
pub type Result<T> = std::result::Result<T, TransitionTransactionError>;
|
||||||
@@ -584,7 +585,8 @@ pub(crate) async fn save_transition_transaction_record(
|
|||||||
let object =
|
let object =
|
||||||
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
|
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
|
||||||
let data = transaction.encode().map_err(transition_transaction_store_error)?;
|
let data = transaction.encode().map_err(transition_transaction_store_error)?;
|
||||||
config_boundary::save_config(api, &object, data).await
|
config_boundary::save_config(api.clone(), &object, data.clone()).await?;
|
||||||
|
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn load_transition_transaction_record(
|
pub(crate) async fn load_transition_transaction_record(
|
||||||
@@ -596,8 +598,14 @@ pub(crate) async fn load_transition_transaction_record(
|
|||||||
TransitionTransaction::decode(transaction_id, &data).map_err(transition_transaction_store_error)
|
TransitionTransaction::decode(transaction_id, &data).map_err(transition_transaction_store_error)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn delete_transition_transaction_record(api: Arc<ECStore>, transaction_id: Uuid) -> EcstoreResult<()> {
|
pub(crate) async fn delete_transition_transaction_record(
|
||||||
let object = transition_transaction_record_object_name(transaction_id).map_err(transition_transaction_store_error)?;
|
api: Arc<ECStore>,
|
||||||
|
transaction: &TransitionTransaction,
|
||||||
|
) -> EcstoreResult<()> {
|
||||||
|
let object =
|
||||||
|
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
|
||||||
|
let data = transaction.encode().map_err(transition_transaction_store_error)?;
|
||||||
|
api.record_durable_ilm_decommission_terminal(&object, &data).await?;
|
||||||
match config_boundary::delete_config(api, &object).await {
|
match config_boundary::delete_config(api, &object).await {
|
||||||
Ok(()) | Err(Error::ConfigNotFound) => Ok(()),
|
Ok(()) | Err(Error::ConfigNotFound) => Ok(()),
|
||||||
Err(err) => Err(err),
|
Err(err) => Err(err),
|
||||||
@@ -813,7 +821,7 @@ pub async fn finalize_missing_transition_transaction_for_operator(
|
|||||||
if probe != TransitionOperatorProbe::Missing {
|
if probe != TransitionOperatorProbe::Missing {
|
||||||
return Err(TransitionOperatorError::CandidateNotMissing(probe));
|
return Err(TransitionOperatorError::CandidateNotMissing(probe));
|
||||||
}
|
}
|
||||||
delete_transition_transaction_record(api, transaction_id)
|
delete_transition_transaction_record(api, &transaction)
|
||||||
.await
|
.await
|
||||||
.map_err(TransitionOperatorError::Store)
|
.map_err(TransitionOperatorError::Store)
|
||||||
}
|
}
|
||||||
@@ -849,22 +857,22 @@ pub async fn process_transition_transaction_record(
|
|||||||
match transaction.state {
|
match transaction.state {
|
||||||
TransitionTransactionState::Uploaded => {
|
TransitionTransactionState::Uploaded => {
|
||||||
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||||
}
|
}
|
||||||
TransitionTransactionState::CleanupPending => match local_commit_matches_transaction(api.clone(), transaction).await {
|
TransitionTransactionState::CleanupPending => match local_commit_matches_transaction(api.clone(), transaction).await {
|
||||||
Ok(true) => {
|
Ok(true) => {
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||||
}
|
}
|
||||||
Ok(false) => {
|
Ok(false) => {
|
||||||
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||||
}
|
}
|
||||||
Err(err) if transition_source_is_missing(&err) => {
|
Err(err) if transition_source_is_missing(&err) => {
|
||||||
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||||
}
|
}
|
||||||
Err(err) => Err(err),
|
Err(err) => Err(err),
|
||||||
@@ -872,7 +880,7 @@ pub async fn process_transition_transaction_record(
|
|||||||
TransitionTransactionState::LocalCommitStarted => {
|
TransitionTransactionState::LocalCommitStarted => {
|
||||||
match local_commit_matches_transaction(api.clone(), transaction).await {
|
match local_commit_matches_transaction(api.clone(), transaction).await {
|
||||||
Ok(true) => {
|
Ok(true) => {
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||||
}
|
}
|
||||||
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||||
@@ -881,7 +889,7 @@ pub async fn process_transition_transaction_record(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed => {
|
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed => {
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||||
}
|
}
|
||||||
TransitionTransactionState::UploadOutcomeUnknown => recover_unknown_upload_outcome(api, transaction).await,
|
TransitionTransactionState::UploadOutcomeUnknown => recover_unknown_upload_outcome(api, transaction).await,
|
||||||
@@ -907,7 +915,7 @@ async fn recover_unknown_upload_outcome(
|
|||||||
.map_err(Error::other)?
|
.map_err(Error::other)?
|
||||||
{
|
{
|
||||||
TransitionCandidateProbe::Missing => {
|
TransitionCandidateProbe::Missing => {
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||||
}
|
}
|
||||||
TransitionCandidateProbe::UnversionedPresent => {
|
TransitionCandidateProbe::UnversionedPresent => {
|
||||||
@@ -925,7 +933,7 @@ async fn recover_unknown_upload_outcome(
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(Error::other)?;
|
.map_err(Error::other)?;
|
||||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
delete_transition_transaction_record(api, transaction).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||||
}
|
}
|
||||||
TransitionCandidateProbe::VersionedPresent(version_id) => {
|
TransitionCandidateProbe::VersionedPresent(version_id) => {
|
||||||
@@ -958,7 +966,7 @@ async fn cleanup_recovered_unknown_upload_candidate(
|
|||||||
.map_err(transition_transaction_store_error)?;
|
.map_err(transition_transaction_store_error)?;
|
||||||
save_transition_transaction_record(api.clone(), &cleanup).await?;
|
save_transition_transaction_record(api.clone(), &cleanup).await?;
|
||||||
delete_transition_remote_candidate(api.clone(), &cleanup).await?;
|
delete_transition_remote_candidate(api.clone(), &cleanup).await?;
|
||||||
delete_transition_transaction_record(api, cleanup.transaction_id).await?;
|
delete_transition_transaction_record(api, &cleanup).await?;
|
||||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -406,6 +406,25 @@ where
|
|||||||
Ok(data)
|
Ok(data)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn read_config_limited_preserve_empty<S>(api: Arc<S>, file: &str, max_bytes: usize) -> Result<Vec<u8>>
|
||||||
|
where
|
||||||
|
S: EcstoreObjectIO,
|
||||||
|
{
|
||||||
|
let (data, _obj) = read_config_limited_preserve_empty_with_metadata(api, file, max_bytes).await?;
|
||||||
|
Ok(data)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn read_config_limited_preserve_empty_with_metadata<S>(
|
||||||
|
api: Arc<S>,
|
||||||
|
file: &str,
|
||||||
|
max_bytes: usize,
|
||||||
|
) -> Result<(Vec<u8>, ObjectInfo)>
|
||||||
|
where
|
||||||
|
S: EcstoreObjectIO,
|
||||||
|
{
|
||||||
|
read_config_with_metadata_inner(api, file, &ObjectOptions::default(), true, Some(max_bytes)).await
|
||||||
|
}
|
||||||
|
|
||||||
/// Read an existing config object without treating an empty payload as absent.
|
/// Read an existing config object without treating an empty payload as absent.
|
||||||
/// Callers that validate their own payload format need to distinguish corruption
|
/// Callers that validate their own payload format need to distinguish corruption
|
||||||
/// from `ConfigNotFound`.
|
/// from `ConfigNotFound`.
|
||||||
|
|||||||
+1607
-47
File diff suppressed because it is too large
Load Diff
@@ -37,7 +37,7 @@ use crate::bucket::lifecycle::{
|
|||||||
transition_transaction::{
|
transition_transaction::{
|
||||||
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
|
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
|
||||||
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
|
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
|
||||||
save_transition_transaction_record,
|
load_transition_transaction_record, save_transition_transaction_record,
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
use crate::bucket::quota::reservation;
|
use crate::bucket::quota::reservation;
|
||||||
@@ -4246,7 +4246,12 @@ fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction,
|
|||||||
|
|
||||||
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
|
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
|
||||||
if let Some(api) = api {
|
if let Some(api) = api {
|
||||||
return delete_transition_transaction_record(api.clone(), transaction_id).await;
|
let transaction = match load_transition_transaction_record(api.clone(), transaction_id).await {
|
||||||
|
Ok(transaction) => transaction,
|
||||||
|
Err(Error::ConfigNotFound) => return Ok(()),
|
||||||
|
Err(err) => return Err(err),
|
||||||
|
};
|
||||||
|
return delete_transition_transaction_record(api.clone(), &transaction).await;
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user