mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-20 03:22:18 +00:00
feat(ecstore): attribute manual transition worker failures (#5324)
* feat(ecstore): attribute manual transition worker failures Add manual transition worker failure reason tracking and persistence recovery compatibility for checksum validation. Co-Authored-By: heihutu <heihutu@gmail.com> * chore(ilm): add manual transition diagnostics scripts Co-Authored-By: heihutu <heihutu@gmail.com> * style(ecstore): format manual transition attribution Co-Authored-By: heihutu <heihutu@gmail.com> * fix(ecstore): avoid copying failure reasons via clone Co-Authored-By: heihutu <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -23,13 +23,13 @@ use crate::bucket::lifecycle::lifecycle::{
|
||||
};
|
||||
use crate::bucket::lifecycle::manual_transition_job::{
|
||||
MANUAL_TRANSITION_JOB_RECORD_PREFIX, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission,
|
||||
ManualTransitionScopeAdmissionClaim, ManualTransitionTaskRecord, ManualTransitionWorkerResult,
|
||||
claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
|
||||
ManualTransitionScopeAdmissionClaim, ManualTransitionTaskRecord, ManualTransitionWorkerFailureReason,
|
||||
ManualTransitionWorkerResult, claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
|
||||
load_manual_transition_job_record, load_manual_transition_job_record_with_etag,
|
||||
manual_transition_job_id_from_record_object_name, manual_transition_job_lease_expired,
|
||||
manual_transition_worker_result_task_key, persist_manual_transition_job_progress, reconcile_manual_transition_worker_results,
|
||||
record_manual_transition_worker_result, renew_manual_transition_job_lease, save_manual_transition_job_record_if_current,
|
||||
save_manual_transition_task_if_absent,
|
||||
record_manual_transition_worker_result, record_manual_transition_worker_result_with_reason,
|
||||
renew_manual_transition_job_lease, save_manual_transition_job_record_if_current, save_manual_transition_task_if_absent,
|
||||
};
|
||||
use crate::bucket::lifecycle::replication_sink;
|
||||
use crate::bucket::lifecycle::replication_sink::{
|
||||
@@ -49,7 +49,9 @@ use crate::disk::error::DiskError;
|
||||
use crate::disk::{DeleteOptions, Disk, DiskAPI, RUSTFS_META_BUCKET, RUSTFS_META_MULTIPART_BUCKET, STORAGE_FORMAT_FILE};
|
||||
use crate::error::Error;
|
||||
use crate::error::StorageError;
|
||||
use crate::error::{error_resp_to_object_err, is_err_object_not_found, is_err_version_not_found, is_network_or_host_down};
|
||||
use crate::error::{
|
||||
error_resp_to_object_err, is_err_object_not_found, is_err_read_quorum, is_err_version_not_found, is_network_or_host_down,
|
||||
};
|
||||
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions};
|
||||
use crate::services::tier::{
|
||||
tier::{TierConfigMgr, TierOperationLease, tier_destination_id_from_metadata},
|
||||
@@ -94,7 +96,7 @@ use s3s::dto::{
|
||||
use s3s::header::{X_AMZ_RESTORE, X_AMZ_SERVER_SIDE_ENCRYPTION};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::any::Any;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::collections::{BTreeMap, HashMap, HashSet};
|
||||
use std::env;
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
@@ -1658,11 +1660,12 @@ impl TransitionState {
|
||||
.await
|
||||
{
|
||||
if let (Some(job_id), Some(result_key)) = (task.manual_job_id, task.manual_result_key.as_deref()) {
|
||||
record_manual_transition_worker_result_for_task(
|
||||
record_manual_transition_worker_result_for_task_with_reason(
|
||||
api.clone(),
|
||||
job_id,
|
||||
result_key,
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
manual_transition_worker_failure_reason(&err),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
@@ -1818,6 +1821,78 @@ async fn record_manual_transition_worker_result_for_task(
|
||||
}
|
||||
}
|
||||
|
||||
async fn record_manual_transition_worker_result_for_task_with_reason(
|
||||
api: Arc<ECStore>,
|
||||
job_id: Uuid,
|
||||
result_key: &str,
|
||||
result: ManualTransitionWorkerResult,
|
||||
reason: ManualTransitionWorkerFailureReason,
|
||||
) {
|
||||
if let Err(err) = record_manual_transition_worker_result_with_reason(
|
||||
api,
|
||||
job_id,
|
||||
result_key,
|
||||
result,
|
||||
manual_transition_queue_snapshot(),
|
||||
Some(reason),
|
||||
)
|
||||
.await
|
||||
{
|
||||
warn!(
|
||||
event = EVENT_LIFECYCLE_WORKER_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
job_id = %job_id,
|
||||
error = %err,
|
||||
reason = ?reason,
|
||||
state = "manual_transition_worker_result_failed",
|
||||
"Manual transition worker failed to persist job result"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn manual_transition_worker_failure_reason(err: &Error) -> ManualTransitionWorkerFailureReason {
|
||||
if is_err_object_not_found(err) || is_err_version_not_found(err) {
|
||||
return ManualTransitionWorkerFailureReason::NotFound;
|
||||
}
|
||||
if is_err_permission_denied(err) {
|
||||
return ManualTransitionWorkerFailureReason::PermissionDenied;
|
||||
}
|
||||
if is_err_read_quorum(err) {
|
||||
return ManualTransitionWorkerFailureReason::Quorum;
|
||||
}
|
||||
if is_timeout(err) {
|
||||
return ManualTransitionWorkerFailureReason::Timeout;
|
||||
}
|
||||
if is_slow_down(err) {
|
||||
return ManualTransitionWorkerFailureReason::SlowDown;
|
||||
}
|
||||
if is_network_or_host_down(&err.to_string(), false) {
|
||||
return ManualTransitionWorkerFailureReason::Network;
|
||||
}
|
||||
ManualTransitionWorkerFailureReason::Unknown
|
||||
}
|
||||
|
||||
fn is_err_permission_denied(err: &Error) -> bool {
|
||||
match err {
|
||||
Error::VolumeAccessDenied | Error::FileAccessDenied | Error::PrefixAccessDenied(_, _) | Error::DiskAccessDenied => true,
|
||||
Error::Io(io) => io.kind() == std::io::ErrorKind::PermissionDenied,
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
fn is_timeout(err: &Error) -> bool {
|
||||
match err {
|
||||
Error::Timeout => true,
|
||||
Error::Io(io) => io.kind() == std::io::ErrorKind::TimedOut,
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
fn is_slow_down(err: &Error) -> bool {
|
||||
matches!(err, Error::SlowDown)
|
||||
}
|
||||
|
||||
pub async fn init_background_expiry(api: Arc<ECStore>) {
|
||||
let mut workers = get_env_usize("RUSTFS_MAX_EXPIRY_WORKERS", std::cmp::min(num_cpus::get(), 16));
|
||||
//globalILMConfig.getExpirationWorkers()
|
||||
@@ -3233,6 +3308,8 @@ pub struct ManualTransitionRunReport {
|
||||
#[serde(default, skip_serializing_if = "is_zero_u64")]
|
||||
pub transition_failed: u64,
|
||||
pub tier_failure: u64,
|
||||
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
|
||||
pub tier_failure_by_reason: BTreeMap<ManualTransitionWorkerFailureReason, u64>,
|
||||
pub truncated_by_limit: bool,
|
||||
pub truncated_by_duration: bool,
|
||||
pub cancelled: bool,
|
||||
@@ -3303,12 +3380,18 @@ impl ManualTransitionRunReport {
|
||||
}
|
||||
|
||||
pub fn merge_scan_report_preserving_worker(&mut self, scan_report: &ManualTransitionRunReport) {
|
||||
let mut tier_failure_by_reason = self.tier_failure_by_reason.clone();
|
||||
for (reason, count) in &scan_report.tier_failure_by_reason {
|
||||
let current = tier_failure_by_reason.get(reason).copied().unwrap_or_default();
|
||||
tier_failure_by_reason.insert(*reason, current.max(*count));
|
||||
}
|
||||
let transition_completed = self.transition_completed;
|
||||
let transition_failed = self.transition_failed;
|
||||
*self = scan_report.clone();
|
||||
self.transition_completed = transition_completed;
|
||||
self.transition_failed = transition_failed;
|
||||
self.tier_failure = scan_report.tier_failure.saturating_add(transition_failed);
|
||||
self.tier_failure_by_reason = tier_failure_by_reason;
|
||||
}
|
||||
|
||||
pub fn worker_transition_pending(&self) -> bool {
|
||||
@@ -4658,12 +4741,13 @@ mod tests {
|
||||
lifecycle_delete_all_versions_replication_scan, lifecycle_deleted_object, lifecycle_replication_blocks_action,
|
||||
lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets,
|
||||
manual_transition_duration_elapsed, manual_transition_has_more_after_limit, manual_transition_recovery_progress_sink,
|
||||
manual_transition_version_marker, mark_delete_opts_skip_decommissioned_on_remote_success,
|
||||
merge_stale_multipart_candidate, persist_manual_transition_job_progress, persist_manual_transition_page_checkpoint,
|
||||
recover_manual_transition_job, recover_manual_transition_jobs, replication_state_for_delete,
|
||||
resolve_tier_free_version_recovery_enabled, resolve_transition_queue_capacity, resolve_transition_queue_send_timeout,
|
||||
resolve_transition_worker_count, resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop,
|
||||
select_restore_s3_location, set_lifecycle_observability_observer, set_recovered_free_version_enqueue_observer,
|
||||
manual_transition_version_marker, manual_transition_worker_failure_reason,
|
||||
mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate,
|
||||
persist_manual_transition_job_progress, persist_manual_transition_page_checkpoint, recover_manual_transition_job,
|
||||
recover_manual_transition_jobs, replication_state_for_delete, resolve_tier_free_version_recovery_enabled,
|
||||
resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count,
|
||||
resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location,
|
||||
set_lifecycle_observability_observer, set_recovered_free_version_enqueue_observer,
|
||||
should_defer_date_expiry_for_recent_config_update, should_reuse_lifecycle_delete_replication_state,
|
||||
transitioned_cleanup_tuple, transitioned_object_delete_opts, wait_for_tier_free_version_recovery,
|
||||
};
|
||||
@@ -4676,16 +4760,17 @@ mod tests {
|
||||
use crate::bucket::lifecycle::config_boundary;
|
||||
use crate::bucket::lifecycle::manual_transition_job::{
|
||||
ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim,
|
||||
ManualTransitionTaskRecord, ManualTransitionWorkerResult, ManualTransitionWorkerResultRecord,
|
||||
claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
|
||||
legacy_manual_transition_scope_key, load_manual_transition_job_record, load_manual_transition_scope_admission,
|
||||
ManualTransitionTaskRecord, ManualTransitionWorkerFailureReason, ManualTransitionWorkerResult,
|
||||
ManualTransitionWorkerResultRecord, claim_manual_transition_scope_admission,
|
||||
delete_manual_transition_scope_admission_if_current, legacy_manual_transition_scope_key,
|
||||
load_manual_transition_job_record, load_manual_transition_scope_admission,
|
||||
load_manual_transition_scope_admission_with_etag, load_manual_transition_task_record,
|
||||
manual_transition_scope_record_object_name, manual_transition_worker_result_object_name,
|
||||
manual_transition_worker_result_task_key, reconcile_manual_transition_worker_results,
|
||||
record_manual_transition_worker_result, renew_manual_transition_job_lease, request_manual_transition_job_cancel,
|
||||
save_manual_transition_job_record, save_manual_transition_scope_admission_if_absent,
|
||||
save_manual_transition_scope_admission_if_current, save_manual_transition_task_if_absent,
|
||||
save_manual_transition_worker_result_if_absent,
|
||||
record_manual_transition_worker_result, record_manual_transition_worker_result_with_reason,
|
||||
renew_manual_transition_job_lease, request_manual_transition_job_cancel, save_manual_transition_job_record,
|
||||
save_manual_transition_scope_admission_if_absent, save_manual_transition_scope_admission_if_current,
|
||||
save_manual_transition_task_if_absent, save_manual_transition_worker_result_if_absent,
|
||||
};
|
||||
use crate::bucket::lifecycle::replication_sink::{
|
||||
ReplicateDecision, ReplicateTargetDecision, ReplicationStatusType, VersionPurgeStatusType,
|
||||
@@ -7539,6 +7624,37 @@ mod tests {
|
||||
assert!(record.completed_at_unix_nanos.is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_worker_failure_reason_classifies_known_failures() {
|
||||
assert_eq!(
|
||||
manual_transition_worker_failure_reason(&Error::ObjectNotFound("bucket".to_string(), "obj".to_string())),
|
||||
ManualTransitionWorkerFailureReason::NotFound
|
||||
);
|
||||
assert_eq!(
|
||||
manual_transition_worker_failure_reason(&Error::Io(std::io::Error::new(
|
||||
std::io::ErrorKind::PermissionDenied,
|
||||
"file access denied",
|
||||
))),
|
||||
ManualTransitionWorkerFailureReason::PermissionDenied
|
||||
);
|
||||
assert_eq!(
|
||||
manual_transition_worker_failure_reason(&Error::ErasureReadQuorum),
|
||||
ManualTransitionWorkerFailureReason::Quorum
|
||||
);
|
||||
assert_eq!(
|
||||
manual_transition_worker_failure_reason(&Error::Timeout),
|
||||
ManualTransitionWorkerFailureReason::Timeout
|
||||
);
|
||||
assert_eq!(
|
||||
manual_transition_worker_failure_reason(&Error::SlowDown),
|
||||
ManualTransitionWorkerFailureReason::SlowDown
|
||||
);
|
||||
assert_eq!(
|
||||
manual_transition_worker_failure_reason(&Error::MethodNotAllowed),
|
||||
ManualTransitionWorkerFailureReason::Unknown
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn manual_transition_counts_already_transitioned_object() {
|
||||
let lc = latest_transition_lifecycle();
|
||||
@@ -8471,6 +8587,49 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn manual_transition_worker_result_stores_failure_reason() {
|
||||
let (_paths, ecstore) = setup_test_env().await;
|
||||
let job_id = Uuid::new_v4();
|
||||
let bucket = format!("manual-worker-failure-reason-{}", job_id.simple());
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
let mut record = ManualTransitionJobRecord::new(job_id, &bucket, &options, "worker-owner");
|
||||
record.complete(
|
||||
ManualTransitionRunReport {
|
||||
bucket: bucket.clone(),
|
||||
enqueued: 1,
|
||||
..Default::default()
|
||||
},
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
);
|
||||
save_manual_transition_job_record(ecstore.clone(), &record)
|
||||
.await
|
||||
.expect("worker result job record should save");
|
||||
|
||||
let task_key = manual_transition_worker_result_task_key(&bucket, "logs/fail", None);
|
||||
let final_record = record_manual_transition_worker_result_with_reason(
|
||||
ecstore.clone(),
|
||||
job_id,
|
||||
&task_key,
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
Some(ManualTransitionWorkerFailureReason::Network),
|
||||
)
|
||||
.await
|
||||
.expect("worker result with failure reason should persist");
|
||||
|
||||
assert_eq!(
|
||||
final_record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::Network),
|
||||
Some(&1)
|
||||
);
|
||||
assert_eq!(final_record.report.transition_failed, 1);
|
||||
assert_eq!(final_record.report.tier_failure, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn manual_transition_worker_result_reconcile_applies_marker_and_releases_admission() {
|
||||
|
||||
@@ -12,6 +12,7 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use rustfs_utils::crypto::{hex_sha256, is_sha256_checksum};
|
||||
@@ -194,6 +195,15 @@ impl ManualTransitionJobRecord {
|
||||
}
|
||||
|
||||
pub fn record_worker_result(&mut self, result: ManualTransitionWorkerResult, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||
self.record_worker_result_with_reason(result, queue_snapshot, None);
|
||||
}
|
||||
|
||||
pub fn record_worker_result_with_reason(
|
||||
&mut self,
|
||||
result: ManualTransitionWorkerResult,
|
||||
queue_snapshot: ManualTransitionQueueSnapshot,
|
||||
failure_reason: Option<ManualTransitionWorkerFailureReason>,
|
||||
) {
|
||||
if self.is_terminal() {
|
||||
return;
|
||||
}
|
||||
@@ -204,6 +214,8 @@ impl ManualTransitionJobRecord {
|
||||
ManualTransitionWorkerResult::TierFailure => {
|
||||
self.report.transition_failed = self.report.transition_failed.saturating_add(1);
|
||||
self.report.tier_failure = self.report.tier_failure.saturating_add(1);
|
||||
let reason = failure_reason.unwrap_or(ManualTransitionWorkerFailureReason::Unknown);
|
||||
*self.report.tier_failure_by_reason.entry(reason).or_insert(0) += 1;
|
||||
}
|
||||
}
|
||||
self.queue_snapshot = queue_snapshot;
|
||||
@@ -215,6 +227,7 @@ impl ManualTransitionJobRecord {
|
||||
&mut self,
|
||||
completed: u64,
|
||||
failed: u64,
|
||||
failure_reasons: &BTreeMap<ManualTransitionWorkerFailureReason, u64>,
|
||||
task_queued: u64,
|
||||
queue_snapshot: ManualTransitionQueueSnapshot,
|
||||
) -> bool {
|
||||
@@ -237,11 +250,17 @@ impl ManualTransitionJobRecord {
|
||||
{
|
||||
return false;
|
||||
}
|
||||
let mut scan_tier_failure_by_reason = self.report.tier_failure_by_reason.clone();
|
||||
for (reason, count) in failure_reasons {
|
||||
let current = scan_tier_failure_by_reason.get(reason).copied().unwrap_or_default();
|
||||
scan_tier_failure_by_reason.insert(*reason, current.max(*count));
|
||||
}
|
||||
let scan_tier_failure = self.report.tier_failure.saturating_sub(self.report.transition_failed);
|
||||
self.report.enqueued = enqueued;
|
||||
self.report.transition_completed = transition_completed;
|
||||
self.report.transition_failed = 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.queue_snapshot = queue_snapshot;
|
||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
self.mark_terminal_if_worker_drained();
|
||||
@@ -401,9 +420,24 @@ impl ManualTransitionJobRecord {
|
||||
return Err(ManualTransitionJobError::Corrupt("content checksum is not a sha256 checksum"));
|
||||
}
|
||||
let mut job = persisted.job;
|
||||
let job_bytes = serde_json::to_vec(&job)?;
|
||||
let mut job_bytes = serde_json::to_vec(&job)?;
|
||||
let actual_checksum = hex_sha256(&job_bytes, ToOwned::to_owned);
|
||||
if persisted.content_sha256 != actual_checksum {
|
||||
let checksum_match = if persisted.content_sha256 == actual_checksum {
|
||||
true
|
||||
} else {
|
||||
let mut job_value = serde_json::to_value(&job)?;
|
||||
if let Some(job_obj) = job_value.as_object_mut() {
|
||||
if let Some(report) = job_obj.get_mut("report").and_then(|value| value.as_object_mut()) {
|
||||
report.remove("tier_failure_by_reason");
|
||||
}
|
||||
} else {
|
||||
return Err(ManualTransitionJobError::ChecksumMismatch);
|
||||
}
|
||||
job_bytes = serde_json::to_vec(&job_value)?;
|
||||
let fallback_checksum = hex_sha256(&job_bytes, ToOwned::to_owned);
|
||||
persisted.content_sha256 == fallback_checksum
|
||||
};
|
||||
if !checksum_match {
|
||||
return Err(ManualTransitionJobError::ChecksumMismatch);
|
||||
}
|
||||
if job.job_id != expected_job_id {
|
||||
@@ -449,6 +483,18 @@ pub enum ManualTransitionWorkerResult {
|
||||
TierFailure,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ManualTransitionWorkerFailureReason {
|
||||
Unknown,
|
||||
NotFound,
|
||||
Network,
|
||||
PermissionDenied,
|
||||
Timeout,
|
||||
Quorum,
|
||||
SlowDown,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct ManualTransitionTaskRecord {
|
||||
@@ -554,17 +600,22 @@ struct PersistedManualTransitionTaskRecord {
|
||||
record: ManualTransitionTaskRecord,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
pub struct ManualTransitionWorkerResultStats {
|
||||
pub completed: u64,
|
||||
pub failed: u64,
|
||||
pub tier_failure_by_reason: BTreeMap<ManualTransitionWorkerFailureReason, u64>,
|
||||
}
|
||||
|
||||
impl ManualTransitionWorkerResultStats {
|
||||
fn record(&mut self, result: ManualTransitionWorkerResult) {
|
||||
fn record(&mut self, result: ManualTransitionWorkerResult, failure_reason: Option<ManualTransitionWorkerFailureReason>) {
|
||||
match result {
|
||||
ManualTransitionWorkerResult::Completed => self.completed = self.completed.saturating_add(1),
|
||||
ManualTransitionWorkerResult::TierFailure => self.failed = self.failed.saturating_add(1),
|
||||
ManualTransitionWorkerResult::TierFailure => {
|
||||
self.failed = self.failed.saturating_add(1);
|
||||
let reason = failure_reason.unwrap_or(ManualTransitionWorkerFailureReason::Unknown);
|
||||
*self.tier_failure_by_reason.entry(reason).or_insert(0) += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -576,16 +627,28 @@ pub struct ManualTransitionWorkerResultRecord {
|
||||
pub job_id: Uuid,
|
||||
pub task_key: String,
|
||||
pub result: ManualTransitionWorkerResult,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub failure_reason: Option<ManualTransitionWorkerFailureReason>,
|
||||
pub completed_at_unix_nanos: i128,
|
||||
}
|
||||
|
||||
impl ManualTransitionWorkerResultRecord {
|
||||
pub fn new(job_id: Uuid, task_key: impl Into<String>, result: ManualTransitionWorkerResult) -> Self {
|
||||
Self::new_with_reason(job_id, task_key, result, None)
|
||||
}
|
||||
|
||||
pub fn new_with_reason(
|
||||
job_id: Uuid,
|
||||
task_key: impl Into<String>,
|
||||
result: ManualTransitionWorkerResult,
|
||||
failure_reason: Option<ManualTransitionWorkerFailureReason>,
|
||||
) -> Self {
|
||||
Self {
|
||||
schema: MANUAL_TRANSITION_WORKER_RESULT_SCHEMA.to_string(),
|
||||
job_id,
|
||||
task_key: task_key.into(),
|
||||
result,
|
||||
failure_reason,
|
||||
completed_at_unix_nanos: OffsetDateTime::now_utc().unix_timestamp_nanos(),
|
||||
}
|
||||
}
|
||||
@@ -620,9 +683,23 @@ impl ManualTransitionWorkerResultRecord {
|
||||
));
|
||||
}
|
||||
let record = persisted.record;
|
||||
let record_bytes = serde_json::to_vec(&record)?;
|
||||
let mut record_bytes = serde_json::to_vec(&record)?;
|
||||
let actual_checksum = hex_sha256(&record_bytes, ToOwned::to_owned);
|
||||
if persisted.content_sha256 != actual_checksum {
|
||||
let checksum_match = if persisted.content_sha256 == actual_checksum {
|
||||
true
|
||||
} else {
|
||||
if let Some(record_value) = serde_json::to_value(&record)?.as_object_mut() {
|
||||
if record.failure_reason.is_none() {
|
||||
record_value.remove("failure_reason");
|
||||
}
|
||||
record_bytes = serde_json::to_vec(record_value)?;
|
||||
} else {
|
||||
return Err(ManualTransitionJobError::ChecksumMismatch);
|
||||
}
|
||||
let fallback_checksum = hex_sha256(&record_bytes, ToOwned::to_owned);
|
||||
persisted.content_sha256 == fallback_checksum
|
||||
};
|
||||
if !checksum_match {
|
||||
return Err(ManualTransitionJobError::ChecksumMismatch);
|
||||
}
|
||||
if record.job_id != expected_job_id {
|
||||
@@ -1120,7 +1197,7 @@ async fn scan_manual_transition_worker_result_journal(
|
||||
Ok(result) => result,
|
||||
Err(err) => return Ok(ManualTransitionWorkerResultJournal::Corrupt(err.to_string())),
|
||||
};
|
||||
stats.record(result.result);
|
||||
stats.record(result.result, result.failure_reason);
|
||||
}
|
||||
if !page.is_truncated {
|
||||
return Ok(ManualTransitionWorkerResultJournal::Stats(stats));
|
||||
@@ -1151,7 +1228,13 @@ pub async fn reconcile_manual_transition_worker_results(
|
||||
};
|
||||
for _ in 0..4 {
|
||||
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
|
||||
let changed = record.apply_worker_result_counts(stats.completed, stats.failed, task_stats.queued, queue_snapshot);
|
||||
let changed = record.apply_worker_result_counts(
|
||||
stats.completed,
|
||||
stats.failed,
|
||||
&stats.tier_failure_by_reason,
|
||||
task_stats.queued,
|
||||
queue_snapshot,
|
||||
);
|
||||
if !changed {
|
||||
return Ok(record);
|
||||
}
|
||||
@@ -1446,7 +1529,18 @@ pub async fn record_manual_transition_worker_result(
|
||||
result: ManualTransitionWorkerResult,
|
||||
queue_snapshot: ManualTransitionQueueSnapshot,
|
||||
) -> EcstoreResult<ManualTransitionJobRecord> {
|
||||
let result_record = ManualTransitionWorkerResultRecord::new(job_id, task_key, result);
|
||||
record_manual_transition_worker_result_with_reason(api, job_id, task_key, result, queue_snapshot, None).await
|
||||
}
|
||||
|
||||
pub async fn record_manual_transition_worker_result_with_reason(
|
||||
api: Arc<ECStore>,
|
||||
job_id: Uuid,
|
||||
task_key: &str,
|
||||
result: ManualTransitionWorkerResult,
|
||||
queue_snapshot: ManualTransitionQueueSnapshot,
|
||||
failure_reason: Option<ManualTransitionWorkerFailureReason>,
|
||||
) -> EcstoreResult<ManualTransitionJobRecord> {
|
||||
let result_record = ManualTransitionWorkerResultRecord::new_with_reason(job_id, task_key, result, failure_reason);
|
||||
if !save_manual_transition_worker_result_if_absent(api.clone(), &result_record).await? {
|
||||
return load_manual_transition_job_record(api, job_id).await;
|
||||
}
|
||||
@@ -1456,7 +1550,7 @@ pub async fn record_manual_transition_worker_result(
|
||||
if record.is_terminal() {
|
||||
return Ok(record);
|
||||
}
|
||||
record.record_worker_result(result, queue_snapshot);
|
||||
record.record_worker_result_with_reason(result, queue_snapshot, failure_reason);
|
||||
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
|
||||
Ok(()) => {
|
||||
if record.is_terminal() {
|
||||
@@ -1652,6 +1746,52 @@ mod tests {
|
||||
assert_eq!(record.report.tier_failure, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_record_records_worker_failure_reasons() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
|
||||
|
||||
record.record_worker_result_with_reason(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
Some(ManualTransitionWorkerFailureReason::NotFound),
|
||||
);
|
||||
record.record_worker_result_with_reason(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
Some(ManualTransitionWorkerFailureReason::Network),
|
||||
);
|
||||
record.record_worker_result_with_reason(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
None,
|
||||
);
|
||||
|
||||
assert_eq!(record.report.transition_failed, 3);
|
||||
assert_eq!(record.report.tier_failure, 3);
|
||||
assert_eq!(
|
||||
record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::NotFound),
|
||||
Some(&1)
|
||||
);
|
||||
assert_eq!(
|
||||
record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::Network),
|
||||
Some(&1)
|
||||
);
|
||||
assert_eq!(
|
||||
record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::Unknown),
|
||||
Some(&1)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_cancel_recovery_finishes_after_worker_drain() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
@@ -1709,6 +1849,67 @@ mod tests {
|
||||
assert_eq!(record.report.tier_failure, 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_scan_progress_preserves_worker_failure_reasons() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
|
||||
record.record_worker_result_with_reason(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
Some(ManualTransitionWorkerFailureReason::NotFound),
|
||||
);
|
||||
|
||||
record.update_running_progress(
|
||||
ManualTransitionRunReport {
|
||||
bucket: "bucket".to_string(),
|
||||
scanned: 3,
|
||||
enqueued: 1,
|
||||
tier_failure: 2,
|
||||
..Default::default()
|
||||
},
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::NotFound),
|
||||
Some(&1)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_apply_worker_result_counts_preserves_existing_failure_reasons() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
|
||||
record.record_worker_result_with_reason(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
ManualTransitionQueueSnapshot::default(),
|
||||
Some(ManualTransitionWorkerFailureReason::NotFound),
|
||||
);
|
||||
|
||||
let mut failure_reasons = BTreeMap::new();
|
||||
failure_reasons.insert(ManualTransitionWorkerFailureReason::PermissionDenied, 2);
|
||||
|
||||
record.apply_worker_result_counts(0, 1, &failure_reasons, 1, ManualTransitionQueueSnapshot::default());
|
||||
|
||||
assert_eq!(
|
||||
record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::NotFound),
|
||||
Some(&1)
|
||||
);
|
||||
assert_eq!(
|
||||
record
|
||||
.report
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::PermissionDenied),
|
||||
Some(&2)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_unknown_checkpoint_persists_counters_and_cursor() {
|
||||
let options = ManualTransitionRunOptions {
|
||||
@@ -1955,6 +2156,62 @@ mod tests {
|
||||
assert!(matches!(err, ManualTransitionJobError::ChecksumMismatch));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_worker_result_record_round_trips_without_failure_reason() {
|
||||
let job_id = Uuid::new_v4();
|
||||
let task_key = manual_transition_worker_result_task_key("bucket", "logs/a", None);
|
||||
let record = ManualTransitionWorkerResultRecord::new_with_reason(
|
||||
job_id,
|
||||
&task_key,
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
Some(ManualTransitionWorkerFailureReason::Network),
|
||||
);
|
||||
let encoded = record.encode().expect("worker result record should encode");
|
||||
let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded worker result should be json");
|
||||
let record_obj = value["record"].as_object_mut().expect("record object should be present");
|
||||
record_obj.remove("failure_reason");
|
||||
let record_bytes = serde_json::to_vec(&value["record"]).expect("record without reason should encode");
|
||||
value["content_sha256"] = serde_json::Value::String(hex_sha256(&record_bytes, ToOwned::to_owned));
|
||||
let stripped = serde_json::to_vec(&value).expect("stripped worker result should encode");
|
||||
|
||||
let decoded = ManualTransitionWorkerResultRecord::decode(job_id, &task_key, &stripped)
|
||||
.expect("worker result without reason should decode");
|
||||
|
||||
assert_eq!(decoded.failure_reason, None);
|
||||
assert_eq!(decoded.result, ManualTransitionWorkerResult::TierFailure);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_worker_result_stats_aggregates_reasons() {
|
||||
let mut stats = ManualTransitionWorkerResultStats::default();
|
||||
|
||||
stats.record(ManualTransitionWorkerResult::Completed, None);
|
||||
stats.record(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
Some(ManualTransitionWorkerFailureReason::PermissionDenied),
|
||||
);
|
||||
stats.record(
|
||||
ManualTransitionWorkerResult::TierFailure,
|
||||
Some(ManualTransitionWorkerFailureReason::PermissionDenied),
|
||||
);
|
||||
stats.record(ManualTransitionWorkerResult::TierFailure, None);
|
||||
|
||||
assert_eq!(stats.completed, 1);
|
||||
assert_eq!(stats.failed, 3);
|
||||
assert_eq!(
|
||||
stats
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::PermissionDenied),
|
||||
Some(&2)
|
||||
);
|
||||
assert_eq!(
|
||||
stats
|
||||
.tier_failure_by_reason
|
||||
.get(&ManualTransitionWorkerFailureReason::Unknown),
|
||||
Some(&1)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_task_record_rejects_checksum_drift() {
|
||||
let job_id = Uuid::new_v4();
|
||||
@@ -1971,6 +2228,24 @@ mod tests {
|
||||
assert!(matches!(err, ManualTransitionJobError::ChecksumMismatch));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_record_round_trips_legacy_report_without_failure_reason_map() {
|
||||
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");
|
||||
let report = value["job"]["report"].as_object_mut().expect("report should be object");
|
||||
report.remove("tier_failure_by_reason");
|
||||
let record_bytes = serde_json::to_vec(&value["job"]).expect("job without report failure map should encode");
|
||||
value["content_sha256"] = serde_json::Value::String(hex_sha256(&record_bytes, ToOwned::to_owned));
|
||||
let legacy = serde_json::to_vec(&value).expect("job without report failure map should encode");
|
||||
|
||||
let decoded = ManualTransitionJobRecord::decode(record.job_id, &legacy).expect("legacy job report should decode");
|
||||
|
||||
assert_eq!(decoded.state, ManualTransitionJobState::Running);
|
||||
assert!(decoded.report.tier_failure_by_reason.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_record_rejects_unknown_report_fields() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
|
||||
Reference in New Issue
Block a user