mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-05 19:55:37 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4567b6b5e4 | |||
| dff2142408 | |||
| 7b2563158b | |||
| 29086ab7bc |
@@ -20,10 +20,9 @@
|
||||
//! journal (`count_requests`) carries the assertion in every one of them.
|
||||
|
||||
use super::common::{BoxError, OdmTestEnv, RawResponse, SeedObject, start_configured_env};
|
||||
use crate::fake_s3_target::{FaultAction, Operation};
|
||||
use crate::fake_s3_target::Operation;
|
||||
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
|
||||
use bytes::Bytes;
|
||||
use futures::{StreamExt, TryStreamExt};
|
||||
use std::time::Duration;
|
||||
|
||||
type TestResult = Result<(), BoxError>;
|
||||
@@ -146,38 +145,14 @@ async fn test_odm_range_burst_overflows_the_pull_queue_without_failing_clients()
|
||||
.await?;
|
||||
|
||||
let body = payload(128 * 1024);
|
||||
let blocker = "queue/blocker.bin";
|
||||
env.seed_source(SOURCE_BUCKET, &[SeedObject::new(blocker, body.clone())]);
|
||||
// The one-chunk range completes immediately; its full background pull
|
||||
// occupies the only slot while the remaining requests fill the queue.
|
||||
env.source.inject_for_key(
|
||||
Operation::GetObject,
|
||||
blocker,
|
||||
FaultAction::SlowSendBody {
|
||||
chunk_bytes: 1024,
|
||||
delay: Duration::from_millis(100),
|
||||
},
|
||||
2,
|
||||
);
|
||||
let response = env
|
||||
.raw_object_request(http::Method::GET, bucket, blocker, &[("range", "bytes=0-1023")])
|
||||
.await?;
|
||||
assert_eq!(response.status, 206);
|
||||
assert_eq!(response.body, body.slice(0..1024));
|
||||
env.wait_for_status_counter(bucket, "/inflight_pulls", 1, SETTLE).await?;
|
||||
|
||||
let keys: Vec<String> = (0..REQUESTS).map(|index| format!("queue/object-{index:03}.bin")).collect();
|
||||
let seeds: Vec<SeedObject> = keys.iter().map(|key| SeedObject::new(key.clone(), body.clone())).collect();
|
||||
env.seed_source(SOURCE_BUCKET, &seeds);
|
||||
|
||||
// Bound source connections below the fixture's limit while still
|
||||
// submitting all 100 requests to the eight-slot background queue.
|
||||
let responses: Vec<RawResponse> = futures::stream::iter(
|
||||
let responses: Vec<RawResponse> = futures::future::try_join_all(
|
||||
keys.iter()
|
||||
.map(|key| env.raw_object_request(http::Method::GET, bucket, key, &[("range", "bytes=0-1023")])),
|
||||
)
|
||||
.buffered(16)
|
||||
.try_collect()
|
||||
.await?;
|
||||
for (key, response) in keys.iter().zip(&responses) {
|
||||
assert_eq!(response.status, 206, "{key}: {}", String::from_utf8_lossy(&response.body));
|
||||
@@ -193,15 +168,6 @@ async fn test_odm_range_burst_overflows_the_pull_queue_without_failing_clients()
|
||||
.wait_for_status_counter(bucket, "/counters/pull_failures_total/queue_full", 1, SETTLE)
|
||||
.await?;
|
||||
assert!(queue_full > 0, "a 100-deep burst must overflow an 8-slot queue");
|
||||
let queue_full = usize::try_from(queue_full)?;
|
||||
assert!(queue_full <= REQUESTS);
|
||||
env.wait_for_status_counter(
|
||||
bucket,
|
||||
"/counters/pulled_objects_total/background",
|
||||
u64::try_from(REQUESTS + 1 - queue_full)?,
|
||||
SETTLE,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let ranged_reads: usize = keys.iter().map(|key| source_get_count(&env, key)).sum();
|
||||
assert!(
|
||||
@@ -209,6 +175,9 @@ async fn test_odm_range_burst_overflows_the_pull_queue_without_failing_clients()
|
||||
"every reader is served from the source: {ranged_reads} GETs for {REQUESTS} readers"
|
||||
);
|
||||
let dropped = keys.iter().filter(|key| source_get_count(&env, key) == 1).count();
|
||||
assert_eq!(dropped, queue_full, "only overflowed keys remain without a background GET");
|
||||
assert!(
|
||||
dropped > 0,
|
||||
"the overflowed keys are the ones with no backfill GET, but every key got one"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -265,13 +265,16 @@ async fn list_through_rejects_a_tampered_continuation_token() -> TestResult {
|
||||
let decoded = String::from_utf8(base64_simd::STANDARD.decode_to_vec(token.as_bytes())?)?;
|
||||
assert!(decoded.contains("\"t\":\"odm-list\""), "the merged token is an envelope: {decoded}");
|
||||
|
||||
let tampered = base64_simd::STANDARD.encode_to_string(decoded.replace("\"v\":1", "\"v\":3").as_bytes());
|
||||
assert_ne!(tampered, token, "the test must change the token version");
|
||||
let query = serde_urlencoded::to_string([("continuation-token", tampered.as_str())])?;
|
||||
let rejected = env.raw_list_objects_v2(bucket, &query).await?;
|
||||
let error_body = String::from_utf8_lossy(&rejected.body);
|
||||
assert_eq!(rejected.status, 400, "a bumped token version is a client error: {}", error_body);
|
||||
assert!(error_body.contains("<Code>InvalidArgument</Code>"), "{error_body}");
|
||||
let tampered = base64_simd::STANDARD.encode_to_string(decoded.replace("\"v\":1", "\"v\":2").as_bytes());
|
||||
let rejected = env
|
||||
.raw_list_objects_v2(bucket, &format!("continuation-token={tampered}"))
|
||||
.await?;
|
||||
assert_eq!(
|
||||
rejected.status,
|
||||
400,
|
||||
"a bumped token version is a client error: {}",
|
||||
String::from_utf8_lossy(&rejected.body)
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -32,7 +32,7 @@ pub mod bucket {
|
||||
pub mod bucket_target_sys {
|
||||
pub use crate::bucket::bucket_target_sys::{
|
||||
AdvancedPutOptions, BucketTargetError, BucketTargetSys, PutObjectOptions, RemoveObjectOptions, S3ClientError,
|
||||
SsecPassthroughCapability, TargetClient, append_version_id_query,
|
||||
SsecPassthroughCapability, TargetClient, UnreadableTargetsPolicy, append_version_id_query,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -69,13 +69,6 @@ pub mod bucket {
|
||||
};
|
||||
}
|
||||
|
||||
pub mod recovery_control {
|
||||
pub use crate::bucket::lifecycle::recovery_control::{
|
||||
IlmRecoveryClassification, IlmRecoveryControlPage, IlmRecoveryControlView, IlmRecoveryProtocol,
|
||||
inspect_recovery_control, list_recovery_controls,
|
||||
};
|
||||
}
|
||||
|
||||
pub mod transition_transaction {
|
||||
pub use crate::bucket::lifecycle::transition_transaction::{
|
||||
TransitionOperatorDeleteResult, TransitionOperatorError, TransitionOperatorProbe, TransitionOperatorStatus,
|
||||
|
||||
@@ -369,6 +369,26 @@ struct SsecPassthroughRecord {
|
||||
recorded_at: Instant,
|
||||
}
|
||||
|
||||
/// What a target write does when the bucket's persisted target set exists but
|
||||
/// cannot be decoded.
|
||||
///
|
||||
/// `docs/architecture/remote-credential-sealing-adr.md` forbids rewriting a
|
||||
/// configuration that could not be fully read, because re-serializing a
|
||||
/// partial in-memory view is the one mechanism by which a configured target
|
||||
/// really disappears. That rule guards against an *unintentional* overwrite,
|
||||
/// so an operator who names the hazard keeps a repair path
|
||||
/// (rustfs/backlog#2309); everything that does not name it stays refused.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
|
||||
pub enum UnreadableTargetsPolicy {
|
||||
/// Refuse the write with [`BucketTargetError::BucketRemoteTargetsUnreadable`].
|
||||
#[default]
|
||||
FailClosed,
|
||||
/// Discard the unreadable set; the target being written becomes the whole
|
||||
/// configuration. Reachable only from an admin request that asked for it
|
||||
/// explicitly, and audited by the caller.
|
||||
Replace,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
pub struct BucketTargetSys {
|
||||
pub arn_remotes_map: Arc<RwLock<HashMap<String, ArnTarget>>>,
|
||||
@@ -791,20 +811,45 @@ impl BucketTargetSys {
|
||||
bucket: &str,
|
||||
target: &BucketTarget,
|
||||
update: bool,
|
||||
unreadable_policy: UnreadableTargetsPolicy,
|
||||
) -> Result<BucketTargets, BucketTargetError> {
|
||||
self.validate_target(bucket, target).await?;
|
||||
|
||||
let mut bucket_targets = match self.list_bucket_targets(bucket).await {
|
||||
Ok(targets) => targets,
|
||||
Err(BucketTargetError::BucketRemoteTargetNotFound { .. }) => BucketTargets::default(),
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
let mut bucket_targets = self.targets_base_for_write(bucket, unreadable_policy).await?;
|
||||
|
||||
Self::upsert_target_entry(&mut bucket_targets.targets, target, update)?;
|
||||
|
||||
Ok(bucket_targets)
|
||||
}
|
||||
|
||||
/// The persisted target set a write merges into.
|
||||
///
|
||||
/// An absent configuration starts from the empty set. An unreadable one is
|
||||
/// refused, because re-serializing a partial view of a set this node could
|
||||
/// not decode is how a configured target disappears for good — unless the
|
||||
/// caller carries the operator's explicit
|
||||
/// [`UnreadableTargetsPolicy::Replace`] opt-in, which discards it
|
||||
/// deliberately (rustfs/backlog#2309).
|
||||
async fn targets_base_for_write(
|
||||
&self,
|
||||
bucket: &str,
|
||||
unreadable_policy: UnreadableTargetsPolicy,
|
||||
) -> Result<BucketTargets, BucketTargetError> {
|
||||
match self.list_bucket_targets(bucket).await {
|
||||
Ok(targets) => Ok(targets),
|
||||
Err(BucketTargetError::BucketRemoteTargetNotFound { .. }) => Ok(BucketTargets::default()),
|
||||
// The opt-in discards only a set this node genuinely cannot read.
|
||||
// A readable set still merges through the arm above, so the policy
|
||||
// can never drop a target that was visible here.
|
||||
Err(BucketTargetError::BucketRemoteTargetsUnreadable { .. })
|
||||
if unreadable_policy == UnreadableTargetsPolicy::Replace =>
|
||||
{
|
||||
Ok(BucketTargets::default())
|
||||
}
|
||||
Err(err) => Err(err),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn validate_target(&self, bucket: &str, target: &BucketTarget) -> Result<(), BucketTargetError> {
|
||||
if !target.target_type.is_valid() {
|
||||
return Err(BucketTargetError::BucketRemoteArnTypeInvalid {
|
||||
@@ -4313,4 +4358,75 @@ mod tests {
|
||||
let window = LastMinuteLatency::new();
|
||||
assert_eq!(window.get_total().avg, Duration::from_secs(0));
|
||||
}
|
||||
|
||||
fn repair_target(bucket: &str, id: &str) -> BucketTarget {
|
||||
BucketTarget {
|
||||
source_bucket: bucket.to_string(),
|
||||
endpoint: "remote.example.com".to_string(),
|
||||
target_bucket: "remote".to_string(),
|
||||
arn: format!("arn:rustfs:replication:us-east-1:{bucket}:{id}"),
|
||||
target_type: BucketTargetType::ReplicationService,
|
||||
region: "us-east-1".to_string(),
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
/// rustfs/backlog#2309: after rustfs/rustfs#7172 an undecodable
|
||||
/// `bucket-targets.json` left the bucket with no API repair path at all.
|
||||
/// The refusal is the default and stays the default; the operator's
|
||||
/// explicit opt-in is the only thing that discards the set, and it starts
|
||||
/// the replacement from empty rather than from a partial view of bytes
|
||||
/// this node never decoded.
|
||||
#[tokio::test]
|
||||
async fn an_unreadable_target_set_is_replaced_only_with_the_explicit_opt_in() {
|
||||
let sys = BucketTargetSys::default();
|
||||
let bucket = "targets-repair-opt-in";
|
||||
sys.mark_targets_unreadable(bucket).await;
|
||||
|
||||
assert!(
|
||||
matches!(
|
||||
sys.targets_base_for_write(bucket, UnreadableTargetsPolicy::FailClosed).await,
|
||||
Err(BucketTargetError::BucketRemoteTargetsUnreadable { .. })
|
||||
),
|
||||
"without the opt-in an unreadable target set must still refuse the write"
|
||||
);
|
||||
assert_eq!(
|
||||
UnreadableTargetsPolicy::default(),
|
||||
UnreadableTargetsPolicy::FailClosed,
|
||||
"a caller that says nothing must get the refusal"
|
||||
);
|
||||
|
||||
let base = sys
|
||||
.targets_base_for_write(bucket, UnreadableTargetsPolicy::Replace)
|
||||
.await
|
||||
.expect("the explicit opt-in must let an operator replace an unreadable set");
|
||||
assert!(
|
||||
base.is_empty(),
|
||||
"the replacement must start from an empty set, never from a partial decode"
|
||||
);
|
||||
}
|
||||
|
||||
/// The opt-in is not a wipe switch. On a set this node can read, both
|
||||
/// policies take the same merge path, so a stray `replace-unreadable=true`
|
||||
/// cannot drop a visible target — which is what makes the flag safe to
|
||||
/// repeat in an operator's repair script.
|
||||
#[tokio::test]
|
||||
async fn the_opt_in_never_discards_a_readable_target_set() {
|
||||
let sys = BucketTargetSys::default();
|
||||
let bucket = "targets-repair-readable";
|
||||
let existing = repair_target(bucket, "keep");
|
||||
sys.targets_map
|
||||
.write()
|
||||
.await
|
||||
.insert(bucket.to_string(), vec![existing.clone()]);
|
||||
|
||||
for policy in [UnreadableTargetsPolicy::FailClosed, UnreadableTargetsPolicy::Replace] {
|
||||
let base = sys
|
||||
.targets_base_for_write(bucket, policy)
|
||||
.await
|
||||
.expect("a readable target set must be readable under either policy");
|
||||
assert_eq!(base.targets.len(), 1, "{policy:?} must keep the persisted target");
|
||||
assert_eq!(base.targets[0].arn, existing.arn, "{policy:?} must not rewrite the persisted target");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ use super::{
|
||||
bucket_lifecycle_ops::{
|
||||
ManualTransitionQueueSnapshot, ManualTransitionRunReport, decode_manual_transition_continuation_token,
|
||||
},
|
||||
manual_transition_job, recovery_control, tier_delete_journal, transition_transaction,
|
||||
manual_transition_job, tier_delete_journal, transition_transaction,
|
||||
};
|
||||
use crate::error::{Error, Result};
|
||||
use crate::services::tier::tier_probe_intent;
|
||||
@@ -41,7 +41,6 @@ pub(crate) enum DurableIlmRecordKind {
|
||||
ManualTransitionScope,
|
||||
ManualTransitionTask,
|
||||
ManualTransitionWorkerResult,
|
||||
RecoveryControl,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -106,14 +105,8 @@ pub(crate) const MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE: DurableIlmNamespace
|
||||
max_record_size: manual_transition_job::MAX_MANUAL_TRANSITION_WORKER_RESULT_RECORD_SIZE,
|
||||
kind: DurableIlmRecordKind::ManualTransitionWorkerResult,
|
||||
};
|
||||
pub(crate) const RECOVERY_CONTROL_NAMESPACE: DurableIlmNamespace = DurableIlmNamespace {
|
||||
name: "recovery-control",
|
||||
prefix: recovery_control::ILM_RECOVERY_CONTROL_PREFIX,
|
||||
max_record_size: recovery_control::MAX_ILM_RECOVERY_CONTROL_SIZE,
|
||||
kind: DurableIlmRecordKind::RecoveryControl,
|
||||
};
|
||||
|
||||
pub(crate) const DURABLE_ILM_NAMESPACES: [DurableIlmNamespace; 10] = [
|
||||
pub(crate) const DURABLE_ILM_NAMESPACES: [DurableIlmNamespace; 9] = [
|
||||
TIER_DELETE_JOURNAL_NAMESPACE,
|
||||
TIER_DELETE_JOURNAL_V6_NAMESPACE,
|
||||
TIER_DELETE_DISPATCH_MANIFEST_NAMESPACE,
|
||||
@@ -123,7 +116,6 @@ pub(crate) const DURABLE_ILM_NAMESPACES: [DurableIlmNamespace; 10] = [
|
||||
MANUAL_TRANSITION_SCOPE_NAMESPACE,
|
||||
MANUAL_TRANSITION_TASK_NAMESPACE,
|
||||
MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE,
|
||||
RECOVERY_CONTROL_NAMESPACE,
|
||||
];
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
@@ -249,18 +241,6 @@ pub(crate) enum DurableIlmRecordCheckpoint {
|
||||
ManualTransitionWorkerResult {
|
||||
content_sha256: String,
|
||||
},
|
||||
RecoveryControl {
|
||||
content_sha256: String,
|
||||
identity_sha256: String,
|
||||
source_generation_sha256: String,
|
||||
first_seen_at_unix_nanos: i64,
|
||||
revision: u64,
|
||||
classification: recovery_control::IlmRecoveryClassification,
|
||||
attempt_count: u64,
|
||||
consecutive_failure_count: u32,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
owner_fence_sha256: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
impl DurableIlmRecordCheckpoint {
|
||||
@@ -274,8 +254,7 @@ impl DurableIlmRecordCheckpoint {
|
||||
| Self::ManualTransitionJob { content_sha256, .. }
|
||||
| Self::ManualTransitionScope { content_sha256, .. }
|
||||
| Self::ManualTransitionTask { content_sha256 }
|
||||
| Self::ManualTransitionWorkerResult { content_sha256 }
|
||||
| Self::RecoveryControl { content_sha256, .. } => content_sha256,
|
||||
| Self::ManualTransitionWorkerResult { content_sha256 } => content_sha256,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -549,51 +528,6 @@ impl DurableIlmRecordCheckpoint {
|
||||
..
|
||||
},
|
||||
) => previous_identity == next_identity && next_updated_at > previous_updated_at,
|
||||
(
|
||||
Self::RecoveryControl {
|
||||
identity_sha256: previous_identity,
|
||||
source_generation_sha256: previous_generation,
|
||||
first_seen_at_unix_nanos: previous_first_seen,
|
||||
revision: previous_revision,
|
||||
classification: previous_classification,
|
||||
attempt_count: previous_attempts,
|
||||
consecutive_failure_count: previous_failures,
|
||||
owner_fence_sha256: previous_owner,
|
||||
..
|
||||
},
|
||||
Self::RecoveryControl {
|
||||
identity_sha256: next_identity,
|
||||
source_generation_sha256: next_generation,
|
||||
first_seen_at_unix_nanos: next_first_seen,
|
||||
revision: next_revision,
|
||||
classification: next_classification,
|
||||
attempt_count: next_attempts,
|
||||
consecutive_failure_count: next_failures,
|
||||
owner_fence_sha256: next_owner,
|
||||
..
|
||||
},
|
||||
) => {
|
||||
let adjacent = previous_identity == next_identity
|
||||
&& previous_first_seen == next_first_seen
|
||||
&& previous_revision.checked_add(1) == Some(*next_revision);
|
||||
let claim = next_owner.is_some()
|
||||
&& *previous_classification == recovery_control::IlmRecoveryClassification::Retrying
|
||||
&& *next_classification == recovery_control::IlmRecoveryClassification::Retrying
|
||||
&& previous_attempts.checked_add(1) == Some(*next_attempts)
|
||||
&& previous_failures == next_failures;
|
||||
let source_refresh = previous_owner.is_some()
|
||||
&& previous_owner == next_owner
|
||||
&& *previous_classification == recovery_control::IlmRecoveryClassification::Retrying
|
||||
&& *next_classification == recovery_control::IlmRecoveryClassification::Retrying
|
||||
&& previous_attempts == next_attempts
|
||||
&& previous_failures == next_failures
|
||||
&& previous_generation != next_generation;
|
||||
let completion = previous_owner.is_some()
|
||||
&& next_owner.is_none()
|
||||
&& previous_generation == next_generation
|
||||
&& previous_attempts == next_attempts;
|
||||
adjacent && (claim || source_refresh || completion)
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
|
||||
@@ -619,14 +553,6 @@ impl DurableIlmRecordCheckpoint {
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if let Self::RecoveryControl { classification, .. } = terminal
|
||||
&& !matches!(
|
||||
classification,
|
||||
recovery_control::IlmRecoveryClassification::Terminal | recovery_control::IlmRecoveryClassification::Abandoned
|
||||
)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if self == terminal || self.validate_successor(terminal).is_ok() {
|
||||
return true;
|
||||
}
|
||||
@@ -726,32 +652,6 @@ impl DurableIlmRecordCheckpoint {
|
||||
.is_some_and(|distance| tier_probe_state_reaches(*previous_state, *terminal_state, distance))
|
||||
&& (!previous_remote_version_known || previous_remote_version == terminal_remote_version)
|
||||
}
|
||||
(
|
||||
Self::RecoveryControl {
|
||||
identity_sha256: previous_identity,
|
||||
source_generation_sha256: previous_generation,
|
||||
first_seen_at_unix_nanos: previous_first_seen,
|
||||
revision: previous_revision,
|
||||
attempt_count: previous_attempts,
|
||||
..
|
||||
},
|
||||
Self::RecoveryControl {
|
||||
identity_sha256: terminal_identity,
|
||||
source_generation_sha256: terminal_generation,
|
||||
first_seen_at_unix_nanos: terminal_first_seen,
|
||||
revision: terminal_revision,
|
||||
attempt_count: terminal_attempts,
|
||||
classification:
|
||||
recovery_control::IlmRecoveryClassification::Terminal | recovery_control::IlmRecoveryClassification::Abandoned,
|
||||
..
|
||||
},
|
||||
) => {
|
||||
previous_identity == terminal_identity
|
||||
&& (previous_generation == terminal_generation || terminal_attempts > previous_attempts)
|
||||
&& previous_first_seen == terminal_first_seen
|
||||
&& terminal_revision > previous_revision
|
||||
&& terminal_attempts >= previous_attempts
|
||||
}
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
@@ -1319,35 +1219,6 @@ pub(crate) fn validate_durable_ilm_record(path: &str, data: &[u8]) -> Result<Val
|
||||
},
|
||||
)
|
||||
}
|
||||
DurableIlmRecordKind::RecoveryControl => {
|
||||
let (protocol, control_id) = recovery_control::recovery_control_id_from_record_object_name(path)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
let control =
|
||||
recovery_control::IlmRecoveryControl::decode(&control_id, data).map_err(|err| Error::other(err.to_string()))?;
|
||||
let canonical = recovery_control::recovery_control_record_object_name(protocol, &control_id)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
if canonical != path || control.identity.protocol != protocol {
|
||||
return Err(Error::other("ILM recovery control path is not canonical"));
|
||||
}
|
||||
let identity_sha256 = checkpoint_hash(&control.identity)?;
|
||||
let source_generation_sha256 = checkpoint_hash(&control.observed_source_generation)?;
|
||||
let owner_fence_sha256 = control.owner.as_ref().map(checkpoint_hash).transpose()?;
|
||||
(
|
||||
"control_id",
|
||||
control_id,
|
||||
DurableIlmRecordCheckpoint::RecoveryControl {
|
||||
content_sha256,
|
||||
identity_sha256,
|
||||
source_generation_sha256,
|
||||
first_seen_at_unix_nanos: control.first_seen_at_unix_nanos,
|
||||
revision: control.revision,
|
||||
classification: control.classification,
|
||||
attempt_count: control.attempt_count,
|
||||
consecutive_failure_count: control.consecutive_failure_count,
|
||||
owner_fence_sha256,
|
||||
},
|
||||
)
|
||||
}
|
||||
DurableIlmRecordKind::ManualTransitionJob => {
|
||||
let job_id = manual_transition_job::manual_transition_job_id_from_record_object_name(path)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
@@ -1541,94 +1412,6 @@ mod tests {
|
||||
.checkpoint
|
||||
}
|
||||
|
||||
fn recovery_control_fixture() -> recovery_control::IlmRecoveryControl {
|
||||
let source_path = "ilm/transition-transactions/records/12/34/1234567890abcdef1234567890abcdef.json";
|
||||
let generation = recovery_control::IlmRecoverySourceGeneration::new(
|
||||
transition_transaction::TRANSITION_TRANSACTION_SCHEMA,
|
||||
"source-etag",
|
||||
"a".repeat(64),
|
||||
vec![recovery_control::IlmRecoverySourceCopy {
|
||||
authority: "pool-0/set-0".to_string(),
|
||||
canonical_path: source_path.to_string(),
|
||||
etag: "source-etag".to_string(),
|
||||
encoded_len: 128,
|
||||
content_sha256: "a".repeat(64),
|
||||
}],
|
||||
)
|
||||
.expect("source generation should build");
|
||||
recovery_control::IlmRecoveryControl::new(
|
||||
recovery_control::IlmRecoveryControlIdentity {
|
||||
protocol: recovery_control::IlmRecoveryProtocol::TransitionTransaction,
|
||||
canonical_source_path: source_path.to_string(),
|
||||
stable_operation_identity: "12345678-90ab-cdef-1234-567890abcdef".to_string(),
|
||||
record_class: "transition_transaction_v1".to_string(),
|
||||
},
|
||||
generation,
|
||||
recovery_control::IlmRecoveryClassification::Retrying,
|
||||
1_000_000_000,
|
||||
recovery_control::IlmRecoveryErrorCode::None,
|
||||
)
|
||||
.expect("recovery control should build")
|
||||
}
|
||||
|
||||
fn recovery_control_checkpoint(control: &recovery_control::IlmRecoveryControl) -> DurableIlmRecordCheckpoint {
|
||||
let control_id = control.identity.source_operation_digest().expect("control id should derive");
|
||||
let path = recovery_control::recovery_control_record_object_name(control.identity.protocol, &control_id)
|
||||
.expect("control path should build");
|
||||
let encoded = control.encode().expect("control should encode");
|
||||
let namespace = classify_durable_ilm_record(&path)
|
||||
.expect("recovery control namespace should classify")
|
||||
.expect("recovery control should be durable");
|
||||
assert_eq!(namespace, &RECOVERY_CONTROL_NAMESPACE);
|
||||
validate_durable_ilm_record(&path, &encoded)
|
||||
.expect("recovery control should validate")
|
||||
.checkpoint
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recovery_control_checkpoint_tracks_claim_retry_and_terminal_generations() {
|
||||
let initial_control = recovery_control_fixture();
|
||||
let initial = recovery_control_checkpoint(&initial_control);
|
||||
|
||||
let mut claimed_control = initial_control;
|
||||
let mut advanced_generation = claimed_control.observed_source_generation.clone();
|
||||
advanced_generation.source_schema = "rustfs-transition-transaction-v2".to_string();
|
||||
claimed_control
|
||||
.claim_for_source_generation("node-a", Uuid::new_v4(), 2_000_000_000, 300_000_000_000, advanced_generation)
|
||||
.expect("control should claim");
|
||||
let claimed = recovery_control_checkpoint(&claimed_control);
|
||||
initial.validate_successor(&claimed).expect("claim should advance receipt");
|
||||
|
||||
let mut retry_control = claimed_control;
|
||||
retry_control
|
||||
.record_retryable_failure(3_000_000_000, recovery_control::IlmRecoveryErrorCode::BackendTimeout)
|
||||
.expect("retry should persist");
|
||||
let retry = recovery_control_checkpoint(&retry_control);
|
||||
claimed.validate_successor(&retry).expect("retry should advance receipt");
|
||||
|
||||
let ready_at = retry_control
|
||||
.next_attempt_at_unix_nanos
|
||||
.expect("retry deadline should persist");
|
||||
let mut terminal_control = retry_control;
|
||||
terminal_control
|
||||
.claim("node-b", Uuid::new_v4(), ready_at, 300_000_000_000)
|
||||
.expect("retry should claim");
|
||||
let reclaimed = recovery_control_checkpoint(&terminal_control);
|
||||
retry.validate_successor(&reclaimed).expect("reclaim should advance receipt");
|
||||
terminal_control
|
||||
.finish_attempt(
|
||||
recovery_control::IlmRecoveryClassification::Terminal,
|
||||
recovery_control::IlmRecoveryErrorCode::None,
|
||||
)
|
||||
.expect("control should terminate");
|
||||
let terminal = recovery_control_checkpoint(&terminal_control);
|
||||
reclaimed
|
||||
.validate_successor(&terminal)
|
||||
.expect("terminal state should advance receipt");
|
||||
assert!(initial.is_predecessor_of_terminal(&terminal));
|
||||
assert!(!initial.is_predecessor_of_terminal(&retry));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_probe_intent_checkpoint_tracks_exact_monotonic_generations() {
|
||||
let initial_intent = tier_probe_intent_fixture();
|
||||
|
||||
@@ -24,7 +24,6 @@ pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs, g
|
||||
mod object_handlers_common;
|
||||
mod object_lock_boundary;
|
||||
pub use self::core as lifecycle;
|
||||
pub mod recovery_control;
|
||||
mod replication_sink;
|
||||
pub mod rule;
|
||||
mod runtime_boundary;
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -23,11 +23,6 @@ use uuid::Uuid;
|
||||
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::recovery_control::{
|
||||
IlmRecoveryClassification, IlmRecoveryControl, IlmRecoveryControlIdentity, IlmRecoveryErrorCode, IlmRecoveryProtocol,
|
||||
ObservedIlmRecoveryControl, load_recovery_control, observe_recovery_source, recovery_control_record_object_name,
|
||||
save_recovery_control_if_absent, save_recovery_control_if_current,
|
||||
};
|
||||
use crate::bucket::lifecycle::tier_sweeper::{
|
||||
delete_confirmed_transition_candidate_exact_with_lease_idempotent,
|
||||
delete_object_from_remote_tier_idempotent_with_manager_and_identity,
|
||||
@@ -49,7 +44,6 @@ const EVENT_LIFECYCLE_TRANSITION_TRANSACTION_RECOVERY: &str = "lifecycle_transit
|
||||
pub const DEFAULT_TRANSITION_TRANSACTION_RECOVERY_LIMIT: usize = 1_000;
|
||||
const TRANSITION_TRANSACTION_RECOVERY_INTERVAL: Duration = Duration::from_secs(60);
|
||||
const TRANSITION_TRANSACTION_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
|
||||
const TRANSITION_RECOVERY_CONTROL_LEASE_NANOS: i64 = 15 * 60 * 1_000_000_000;
|
||||
pub const TRANSITION_TRANSACTION_SCHEMA: &str = "rustfs-transition-transaction-v1";
|
||||
pub const TRANSITION_TRANSACTION_PREFIX: &str = "ilm/transition-transactions";
|
||||
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = TRANSITION_TRANSACTION_NAMESPACE.prefix;
|
||||
@@ -743,8 +737,6 @@ pub enum TransitionTransactionRecoveryOutcome {
|
||||
RemoteCandidateDeleted,
|
||||
RecordDeleted,
|
||||
Retained,
|
||||
RetainedAmbiguous(IlmRecoveryErrorCode),
|
||||
OperatorRequired(IlmRecoveryErrorCode),
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -825,80 +817,6 @@ async fn pause_before_transition_recovery_claim(transaction_id: Uuid) {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[derive(Default)]
|
||||
struct TransitionRecoveryTerminalBarrierState {
|
||||
transaction_id: Uuid,
|
||||
arrived: tokio::sync::Notify,
|
||||
release: tokio::sync::Notify,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) struct TransitionRecoveryTerminalBarrier {
|
||||
state: Arc<TransitionRecoveryTerminalBarrierState>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static TRANSITION_RECOVERY_TERMINAL_BARRIER: std::sync::OnceLock<
|
||||
std::sync::Mutex<Option<Arc<TransitionRecoveryTerminalBarrierState>>>,
|
||||
> = std::sync::OnceLock::new();
|
||||
|
||||
#[cfg(test)]
|
||||
impl TransitionRecoveryTerminalBarrier {
|
||||
pub(crate) fn install(transaction_id: Uuid) -> Self {
|
||||
let state = Arc::new(TransitionRecoveryTerminalBarrierState {
|
||||
transaction_id,
|
||||
..Default::default()
|
||||
});
|
||||
let mut slot = TRANSITION_RECOVERY_TERMINAL_BARRIER
|
||||
.get_or_init(|| std::sync::Mutex::new(None))
|
||||
.lock()
|
||||
.expect("transition recovery terminal barrier mutex should not poison");
|
||||
assert!(
|
||||
slot.is_none(),
|
||||
"transition recovery terminal barrier must be installed by one test 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 recovery should persist terminal control before source cleanup");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Drop for TransitionRecoveryTerminalBarrier {
|
||||
fn drop(&mut self) {
|
||||
self.state.release.notify_one();
|
||||
let mut slot = TRANSITION_RECOVERY_TERMINAL_BARRIER
|
||||
.get_or_init(|| std::sync::Mutex::new(None))
|
||||
.lock()
|
||||
.expect("transition recovery terminal barrier mutex should not poison");
|
||||
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
||||
*slot = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn pause_after_transition_recovery_terminal(transaction_id: Uuid) {
|
||||
let barrier = TRANSITION_RECOVERY_TERMINAL_BARRIER
|
||||
.get_or_init(|| std::sync::Mutex::new(None))
|
||||
.lock()
|
||||
.expect("transition recovery terminal barrier mutex should not poison")
|
||||
.as_ref()
|
||||
.filter(|barrier| barrier.transaction_id == transaction_id)
|
||||
.cloned();
|
||||
if let Some(barrier) = barrier {
|
||||
barrier.arrived.notify_one();
|
||||
barrier.release.notified().await;
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum TransitionOperatorProbe {
|
||||
@@ -1102,35 +1020,17 @@ fn transition_transaction_id_from_record_object_name(object: &str) -> Result<Uui
|
||||
let suffix = object
|
||||
.strip_prefix(&prefix)
|
||||
.ok_or(TransitionTransactionError::Corrupt("transaction record path has wrong prefix"))?;
|
||||
let mut parts = suffix.split('/');
|
||||
let shard_a = parts
|
||||
let file_name = suffix
|
||||
.rsplit('/')
|
||||
.next()
|
||||
.ok_or(TransitionTransactionError::Corrupt("transaction record path is incomplete"))?;
|
||||
let shard_b = parts
|
||||
.next()
|
||||
.ok_or(TransitionTransactionError::Corrupt("transaction record path is incomplete"))?;
|
||||
let file_name = parts
|
||||
.next()
|
||||
.ok_or(TransitionTransactionError::Corrupt("transaction record path is incomplete"))?;
|
||||
if parts.next().is_some() {
|
||||
return Err(TransitionTransactionError::Corrupt("transaction record path is not canonical"));
|
||||
}
|
||||
let transaction_key = file_name
|
||||
.strip_suffix(".json")
|
||||
.ok_or(TransitionTransactionError::Corrupt("transaction record path has wrong suffix"))?;
|
||||
if transaction_key.len() != 32
|
||||
|| !transaction_key
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
|
||||
|| shard_a != &transaction_key[..2]
|
||||
|| shard_b != &transaction_key[2..4]
|
||||
{
|
||||
if transaction_key.len() != 32 || !transaction_key.bytes().all(|byte| byte.is_ascii_hexdigit()) {
|
||||
return Err(TransitionTransactionError::Corrupt("transaction record path has invalid transaction id"));
|
||||
}
|
||||
Uuid::parse_str(transaction_key)
|
||||
.ok()
|
||||
.filter(|transaction_id| !transaction_id.is_nil())
|
||||
.ok_or(TransitionTransactionError::Corrupt("transaction record path has invalid uuid"))
|
||||
Uuid::parse_str(transaction_key).map_err(|_| TransitionTransactionError::Corrupt("transaction record path has invalid uuid"))
|
||||
}
|
||||
|
||||
pub async fn process_transition_transaction_record(
|
||||
@@ -1155,27 +1055,6 @@ async fn process_transition_transaction_record_at(
|
||||
) -> EcstoreResult<TransitionTransactionRecoveryOutcome> {
|
||||
let record_name =
|
||||
transition_transaction_record_object_name(observed.transaction_id).map_err(transition_transaction_store_error)?;
|
||||
let now_unix_nanos =
|
||||
i64::try_from(now_unix_nanos).map_err(|_| Error::other("transition transaction recovery timestamp does not fit i64"))?;
|
||||
let recovery_control_identity = transition_recovery_control_identity(observed, &record_name);
|
||||
let recovery_control_id = recovery_control_identity
|
||||
.source_operation_digest()
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
let control_record_name =
|
||||
recovery_control_record_object_name(IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
let control_lock = if transition_state_needs_recovery_control(observed, now_unix_nanos) {
|
||||
Some(
|
||||
api.new_ns_lock(RUSTFS_META_BUCKET, &format!("{control_record_name}.recovery-lock"))
|
||||
.await?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let _control_guard = match &control_lock {
|
||||
Some(lock) => Some(lock.get_write_lock(crate::set_disk::get_lock_acquire_timeout()).await?),
|
||||
None => None,
|
||||
};
|
||||
// The synthetic key avoids nesting the recovery lock with the config
|
||||
// object's own I/O lock. Holding it across the bounded source proof and
|
||||
// remote DELETE elects one destructive recovery worker across nodes.
|
||||
@@ -1194,400 +1073,55 @@ async fn process_transition_transaction_record_at(
|
||||
return Ok(TransitionTransactionRecoveryOutcome::Retained);
|
||||
}
|
||||
|
||||
let mut recovery_control = if transition_state_needs_recovery_control(¤t, now_unix_nanos) {
|
||||
if cleanup_terminal_transition_recovery_control(
|
||||
api.clone(),
|
||||
¤t,
|
||||
&record_name,
|
||||
&recovery_control_identity,
|
||||
&recovery_control_id,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
return Ok(TransitionTransactionRecoveryOutcome::RecordDeleted);
|
||||
}
|
||||
match claim_transition_recovery_control(
|
||||
api.clone(),
|
||||
¤t,
|
||||
&record_name,
|
||||
recovery_control_identity,
|
||||
&recovery_control_id,
|
||||
now_unix_nanos,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
Some(control) => Some(control),
|
||||
None => return Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let recovery = match current.state {
|
||||
match current.state {
|
||||
TransitionTransactionState::Uploaded => {
|
||||
if transition_transaction_ownership_is_active(¤t, i128::from(now_unix_nanos)) {
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained)
|
||||
} else {
|
||||
let mut cleanup = current.clone();
|
||||
cleanup
|
||||
.mark_cleanup_pending(
|
||||
current.fence(),
|
||||
TransitionCleanupProof {
|
||||
transaction_id: current.transaction_id,
|
||||
write_id: current.write_id,
|
||||
remote_object: current.remote_object.clone(),
|
||||
remote_version: current.remote_version.clone(),
|
||||
backend_fingerprint: current.backend_fingerprint,
|
||||
decision: TransitionCleanupDecision::UploadAbortedBeforeLocalCommit,
|
||||
},
|
||||
)
|
||||
.map_err(transition_transaction_store_error)?;
|
||||
#[cfg(test)]
|
||||
pause_before_transition_recovery_claim(current.transaction_id).await;
|
||||
match save_transition_transaction_record_if_current(api.clone(), ¤t, &cleanup).await {
|
||||
Ok(()) => recover_cleanup_pending(api.clone(), &cleanup).await,
|
||||
Err(Error::PreconditionFailed) | Err(Error::ConfigNotFound) => {
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained)
|
||||
}
|
||||
Err(err) => Err(err),
|
||||
}
|
||||
if transition_transaction_ownership_is_active(¤t, now_unix_nanos) {
|
||||
return Ok(TransitionTransactionRecoveryOutcome::Retained);
|
||||
}
|
||||
let mut cleanup = current.clone();
|
||||
cleanup
|
||||
.mark_cleanup_pending(
|
||||
current.fence(),
|
||||
TransitionCleanupProof {
|
||||
transaction_id: current.transaction_id,
|
||||
write_id: current.write_id,
|
||||
remote_object: current.remote_object.clone(),
|
||||
remote_version: current.remote_version.clone(),
|
||||
backend_fingerprint: current.backend_fingerprint,
|
||||
decision: TransitionCleanupDecision::UploadAbortedBeforeLocalCommit,
|
||||
},
|
||||
)
|
||||
.map_err(transition_transaction_store_error)?;
|
||||
#[cfg(test)]
|
||||
pause_before_transition_recovery_claim(current.transaction_id).await;
|
||||
match save_transition_transaction_record_if_current(api.clone(), ¤t, &cleanup).await {
|
||||
Ok(()) => recover_cleanup_pending(api, &cleanup).await,
|
||||
Err(Error::PreconditionFailed) | Err(Error::ConfigNotFound) => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||
Err(err) => Err(err),
|
||||
}
|
||||
}
|
||||
TransitionTransactionState::CleanupPending => recover_cleanup_pending(api.clone(), ¤t).await,
|
||||
TransitionTransactionState::CleanupPending => recover_cleanup_pending(api, ¤t).await,
|
||||
TransitionTransactionState::LocalCommitStarted => match local_commit_matches_transaction(api.clone(), ¤t).await {
|
||||
Ok(true) => Ok(TransitionTransactionRecoveryOutcome::RecordDeleted),
|
||||
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::OperatorRequired(
|
||||
IlmRecoveryErrorCode::LocalCommitAmbiguous,
|
||||
)),
|
||||
Err(err) if transition_source_is_missing(&err) => Ok(TransitionTransactionRecoveryOutcome::OperatorRequired(
|
||||
IlmRecoveryErrorCode::LocalCommitAmbiguous,
|
||||
)),
|
||||
Ok(true) => {
|
||||
delete_transition_transaction_record(api, ¤t).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||
Err(err) if transition_source_is_missing(&err) => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed => {
|
||||
delete_transition_transaction_record(api, ¤t).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
TransitionTransactionState::UploadOutcomeUnknown => {
|
||||
if transition_transaction_ownership_is_active(¤t, i128::from(now_unix_nanos)) {
|
||||
if transition_transaction_ownership_is_active(¤t, now_unix_nanos) {
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained)
|
||||
} else {
|
||||
recover_unknown_upload_outcome(api.clone(), ¤t).await
|
||||
recover_unknown_upload_outcome(api, ¤t).await
|
||||
}
|
||||
}
|
||||
TransitionTransactionState::UploadStarted => {
|
||||
if transition_transaction_ownership_is_active(¤t, i128::from(now_unix_nanos)) {
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained)
|
||||
} else {
|
||||
Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
|
||||
IlmRecoveryErrorCode::RemoteVersionUnknown,
|
||||
))
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(mut control) = recovery_control.take() {
|
||||
let source_to_delete = if matches!(
|
||||
recovery,
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted
|
||||
| TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
) {
|
||||
let refreshed =
|
||||
refresh_transition_recovery_control_source(api.clone(), control, &record_name, current.transaction_id).await?;
|
||||
control = refreshed.0;
|
||||
refreshed.1
|
||||
} else {
|
||||
None
|
||||
};
|
||||
persist_transition_recovery_result(api.clone(), control, &recovery, now_unix_nanos).await?;
|
||||
if let Some(source) = source_to_delete {
|
||||
#[cfg(test)]
|
||||
pause_after_transition_recovery_terminal(source.transaction_id).await;
|
||||
delete_transition_transaction_record(api, &source).await?;
|
||||
}
|
||||
} else if matches!(
|
||||
recovery,
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted | TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
) {
|
||||
delete_transition_transaction_record(api, ¤t).await?;
|
||||
}
|
||||
recovery
|
||||
}
|
||||
|
||||
fn transition_recovery_control_identity(transaction: &TransitionTransaction, record_name: &str) -> IlmRecoveryControlIdentity {
|
||||
IlmRecoveryControlIdentity {
|
||||
protocol: IlmRecoveryProtocol::TransitionTransaction,
|
||||
canonical_source_path: record_name.to_string(),
|
||||
stable_operation_identity: transaction.transaction_id.to_string(),
|
||||
record_class: "transition_transaction_v1".to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn transition_recovery_control_id(transaction: &TransitionTransaction) -> Result<String> {
|
||||
let record_name = transition_transaction_record_object_name(transaction.transaction_id)?;
|
||||
transition_recovery_control_identity(transaction, &record_name)
|
||||
.source_operation_digest()
|
||||
.map_err(|_| TransitionTransactionError::Corrupt("transition recovery control identity is invalid"))
|
||||
}
|
||||
|
||||
fn transition_state_needs_recovery_control(transaction: &TransitionTransaction, now_unix_nanos: i64) -> bool {
|
||||
now_unix_nanos >= transaction.not_after_unix_nanos
|
||||
&& !matches!(
|
||||
transaction.state,
|
||||
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed
|
||||
)
|
||||
}
|
||||
|
||||
async fn cleanup_terminal_transition_recovery_control(
|
||||
api: Arc<ECStore>,
|
||||
transaction: &TransitionTransaction,
|
||||
record_name: &str,
|
||||
identity: &IlmRecoveryControlIdentity,
|
||||
control_id: &str,
|
||||
) -> EcstoreResult<bool> {
|
||||
let observed = match load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await {
|
||||
Ok(observed) => observed,
|
||||
Err(Error::ConfigNotFound) => return Ok(false),
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
if observed.control.classification != IlmRecoveryClassification::Terminal {
|
||||
return Ok(false);
|
||||
}
|
||||
let source = observe_recovery_source(api.clone(), record_name, TRANSITION_TRANSACTION_SCHEMA).await?;
|
||||
let exact_source = source.is_consistent()
|
||||
&& source.generation == observed.control.observed_source_generation
|
||||
&& source.canonical_data.as_deref().is_some_and(|data| {
|
||||
TransitionTransaction::decode(transaction.transaction_id, data).is_ok_and(|decoded| decoded == *transaction)
|
||||
});
|
||||
if observed.control.identity != *identity || !exact_source {
|
||||
return Ok(false);
|
||||
}
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn claim_transition_recovery_control(
|
||||
api: Arc<ECStore>,
|
||||
transaction: &TransitionTransaction,
|
||||
record_name: &str,
|
||||
identity: IlmRecoveryControlIdentity,
|
||||
control_id: &str,
|
||||
now_unix_nanos: i64,
|
||||
) -> EcstoreResult<Option<ObservedIlmRecoveryControl>> {
|
||||
let existing = match load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await {
|
||||
Ok(control) => Some(control),
|
||||
Err(Error::ConfigNotFound) => None,
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
if let Some(observed) = existing.as_ref() {
|
||||
if observed.control.identity != identity {
|
||||
return Ok(None);
|
||||
}
|
||||
if observed
|
||||
.control
|
||||
.owner
|
||||
.as_ref()
|
||||
.is_some_and(|owner| owner.lease_expires_at_unix_nanos <= now_unix_nanos)
|
||||
{
|
||||
let mut expired = observed.control.clone();
|
||||
expired
|
||||
.record_expired_attempt(now_unix_nanos)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
save_recovery_control_if_current(api, observed, &expired).await?;
|
||||
return Ok(None);
|
||||
}
|
||||
if !observed.control.should_attempt_at(now_unix_nanos) {
|
||||
return Ok(None);
|
||||
}
|
||||
}
|
||||
|
||||
let source = match observe_recovery_source(api.clone(), record_name, TRANSITION_TRANSACTION_SCHEMA).await {
|
||||
Ok(source) => source,
|
||||
Err(err) => {
|
||||
if let Some(observed) = existing {
|
||||
persist_transition_recovery_source_failure(api, observed, now_unix_nanos).await?;
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(err);
|
||||
}
|
||||
};
|
||||
let source_matches = source.is_consistent()
|
||||
&& source.canonical_data.as_deref().is_some_and(|data| {
|
||||
TransitionTransaction::decode(transaction.transaction_id, data).is_ok_and(|observed| observed == *transaction)
|
||||
});
|
||||
let source_error = if source_matches {
|
||||
IlmRecoveryErrorCode::None
|
||||
} else if source.canonical_data.is_some() {
|
||||
IlmRecoveryErrorCode::SourceGenerationChanged
|
||||
} else {
|
||||
IlmRecoveryErrorCode::SourceDivergent
|
||||
};
|
||||
|
||||
let mut observed = match existing {
|
||||
Some(control) => control,
|
||||
None => {
|
||||
let candidate = IlmRecoveryControl::new(
|
||||
identity.clone(),
|
||||
source.generation.clone(),
|
||||
if source_matches {
|
||||
IlmRecoveryClassification::Retrying
|
||||
} else {
|
||||
IlmRecoveryClassification::Corrupt
|
||||
},
|
||||
now_unix_nanos,
|
||||
source_error,
|
||||
)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
match save_recovery_control_if_absent(api.clone(), &candidate).await {
|
||||
Ok(()) | Err(Error::PreconditionFailed) => {}
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await?
|
||||
}
|
||||
};
|
||||
if observed.control.identity != identity || !observed.control.should_attempt_at(now_unix_nanos) {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let mut claimed = observed.control.clone();
|
||||
claimed
|
||||
.claim_for_source_generation(
|
||||
api.id.to_string(),
|
||||
Uuid::new_v4(),
|
||||
now_unix_nanos,
|
||||
TRANSITION_RECOVERY_CONTROL_LEASE_NANOS,
|
||||
source.generation,
|
||||
)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
save_recovery_control_if_current(api.clone(), &observed, &claimed).await?;
|
||||
observed = load_recovery_control(api.clone(), IlmRecoveryProtocol::TransitionTransaction, control_id).await?;
|
||||
if observed.control != claimed {
|
||||
return Err(Error::PreconditionFailed);
|
||||
}
|
||||
if !source_matches {
|
||||
let mut corrupt = observed.control.clone();
|
||||
corrupt
|
||||
.finish_attempt(IlmRecoveryClassification::Corrupt, source_error)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
save_recovery_control_if_current(api, &observed, &corrupt).await?;
|
||||
return Ok(None);
|
||||
}
|
||||
Ok(Some(observed))
|
||||
}
|
||||
|
||||
async fn persist_transition_recovery_source_failure(
|
||||
api: Arc<ECStore>,
|
||||
observed: ObservedIlmRecoveryControl,
|
||||
now_unix_nanos: i64,
|
||||
) -> EcstoreResult<()> {
|
||||
let mut claimed = observed.control.clone();
|
||||
claimed
|
||||
.claim(
|
||||
api.id.to_string(),
|
||||
Uuid::new_v4(),
|
||||
now_unix_nanos,
|
||||
TRANSITION_RECOVERY_CONTROL_LEASE_NANOS,
|
||||
)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
save_recovery_control_if_current(api.clone(), &observed, &claimed).await?;
|
||||
let claimed = load_recovery_control(
|
||||
api.clone(),
|
||||
IlmRecoveryProtocol::TransitionTransaction,
|
||||
&claimed
|
||||
.identity
|
||||
.source_operation_digest()
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
)
|
||||
.await?;
|
||||
let mut failed = claimed.control.clone();
|
||||
failed
|
||||
.record_retryable_failure(now_unix_nanos, IlmRecoveryErrorCode::SourceUnavailable)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
save_recovery_control_if_current(api, &claimed, &failed).await
|
||||
}
|
||||
|
||||
async fn refresh_transition_recovery_control_source(
|
||||
api: Arc<ECStore>,
|
||||
mut observed: ObservedIlmRecoveryControl,
|
||||
record_name: &str,
|
||||
transaction_id: Uuid,
|
||||
) -> EcstoreResult<(ObservedIlmRecoveryControl, Option<TransitionTransaction>)> {
|
||||
let transaction = match load_transition_transaction_record(api.clone(), transaction_id).await {
|
||||
Ok(transaction) => transaction,
|
||||
Err(Error::ConfigNotFound) => return Ok((observed, None)),
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
let source = observe_recovery_source(api.clone(), record_name, TRANSITION_TRANSACTION_SCHEMA).await?;
|
||||
let exact_source = source.is_consistent()
|
||||
&& source
|
||||
.canonical_data
|
||||
.as_deref()
|
||||
.is_some_and(|data| TransitionTransaction::decode(transaction_id, data).is_ok_and(|decoded| decoded == transaction));
|
||||
if !exact_source {
|
||||
return Err(Error::PreconditionFailed);
|
||||
}
|
||||
if observed.control.observed_source_generation != source.generation {
|
||||
let mut refreshed = observed.control.clone();
|
||||
refreshed
|
||||
.refresh_owned_source_generation(source.generation)
|
||||
.map_err(|err| Error::other(err.to_string()))?;
|
||||
save_recovery_control_if_current(api.clone(), &observed, &refreshed).await?;
|
||||
observed = load_recovery_control(
|
||||
api,
|
||||
IlmRecoveryProtocol::TransitionTransaction,
|
||||
&refreshed
|
||||
.identity
|
||||
.source_operation_digest()
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
)
|
||||
.await?;
|
||||
if observed.control != refreshed {
|
||||
return Err(Error::PreconditionFailed);
|
||||
}
|
||||
}
|
||||
Ok((observed, Some(transaction)))
|
||||
}
|
||||
|
||||
async fn persist_transition_recovery_result(
|
||||
api: Arc<ECStore>,
|
||||
observed: ObservedIlmRecoveryControl,
|
||||
recovery: &EcstoreResult<TransitionTransactionRecoveryOutcome>,
|
||||
now_unix_nanos: i64,
|
||||
) -> EcstoreResult<()> {
|
||||
let mut next = observed.control.clone();
|
||||
match recovery {
|
||||
Ok(
|
||||
TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted | TransitionTransactionRecoveryOutcome::RecordDeleted,
|
||||
) => next
|
||||
.finish_attempt(IlmRecoveryClassification::Terminal, IlmRecoveryErrorCode::None)
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained) => next
|
||||
.record_retryable_failure(now_unix_nanos, IlmRecoveryErrorCode::SourceGenerationChanged)
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(code)) => next
|
||||
.finish_attempt(IlmRecoveryClassification::RetainedAmbiguous, *code)
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
Ok(TransitionTransactionRecoveryOutcome::OperatorRequired(code)) => next
|
||||
.finish_attempt(IlmRecoveryClassification::OperatorRequired, *code)
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
Err(err) => next
|
||||
.record_retryable_failure(now_unix_nanos, transition_recovery_error_code(err))
|
||||
.map_err(|err| Error::other(err.to_string()))?,
|
||||
}
|
||||
save_recovery_control_if_current(api, &observed, &next).await
|
||||
}
|
||||
|
||||
fn transition_recovery_error_code(err: &Error) -> IlmRecoveryErrorCode {
|
||||
match err {
|
||||
Error::PreconditionFailed => IlmRecoveryErrorCode::CasConflict,
|
||||
Error::ConfigNotFound
|
||||
| Error::FileNotFound
|
||||
| Error::FileVersionNotFound
|
||||
| Error::ObjectNotFound(_, _)
|
||||
| Error::VersionNotFound(_, _, _)
|
||||
| Error::BucketNotFound(_) => IlmRecoveryErrorCode::SourceUnavailable,
|
||||
Error::SlowDown => IlmRecoveryErrorCode::BackendThrottled,
|
||||
_ => IlmRecoveryErrorCode::Unknown,
|
||||
TransitionTransactionState::UploadStarted => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1600,7 +1134,10 @@ async fn recover_cleanup_pending(
|
||||
transaction: &TransitionTransaction,
|
||||
) -> EcstoreResult<TransitionTransactionRecoveryOutcome> {
|
||||
match local_commit_matches_transaction(api.clone(), transaction).await {
|
||||
Ok(true) => Ok(TransitionTransactionRecoveryOutcome::RecordDeleted),
|
||||
Ok(true) => {
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
Ok(false) => delete_unreferenced_transition_candidate(api, transaction).await,
|
||||
Err(err) if transition_source_is_missing(&err) => delete_unreferenced_transition_candidate(api, transaction).await,
|
||||
Err(err) => Err(err),
|
||||
@@ -1620,6 +1157,7 @@ async fn delete_unreferenced_transition_candidate(
|
||||
return Ok(TransitionTransactionRecoveryOutcome::Retained);
|
||||
}
|
||||
delete_transition_remote_candidate(api.clone(), ¤t).await?;
|
||||
delete_transition_transaction_record(api, ¤t).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||
}
|
||||
|
||||
@@ -1640,26 +1178,24 @@ async fn recover_unknown_upload_outcome(
|
||||
.await
|
||||
.map_err(Error::other)?
|
||||
{
|
||||
TransitionCandidateProbe::Missing => Ok(TransitionTransactionRecoveryOutcome::RecordDeleted),
|
||||
TransitionCandidateProbe::Missing => {
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
TransitionCandidateProbe::UnversionedPresent => {
|
||||
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
|
||||
}
|
||||
TransitionCandidateProbe::VersionedPresent(version_id)
|
||||
if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) =>
|
||||
{
|
||||
Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
|
||||
IlmRecoveryErrorCode::RemoteVersionUnknown,
|
||||
))
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained)
|
||||
}
|
||||
TransitionCandidateProbe::VersionedPresent(version_id) => {
|
||||
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await
|
||||
}
|
||||
TransitionCandidateProbe::Ambiguous => Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
|
||||
IlmRecoveryErrorCode::RemoteProbeAmbiguous,
|
||||
)),
|
||||
TransitionCandidateProbe::Unsupported => Ok(TransitionTransactionRecoveryOutcome::RetainedAmbiguous(
|
||||
IlmRecoveryErrorCode::RemoteProbeUnsupported,
|
||||
)),
|
||||
TransitionCandidateProbe::Ambiguous | TransitionCandidateProbe::Unsupported => {
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1787,11 +1323,6 @@ async fn recover_transition_transaction_records_with_now(
|
||||
false,
|
||||
)
|
||||
.await?;
|
||||
if list.is_truncated && list.next_continuation_token.is_none() {
|
||||
return Err(Error::other(
|
||||
"transition transaction recovery returned a truncated page without a continuation marker",
|
||||
));
|
||||
}
|
||||
|
||||
let mut stats = TransitionTransactionRecoveryStats {
|
||||
scanned: 0,
|
||||
@@ -1850,11 +1381,7 @@ async fn recover_transition_transaction_records_with_now(
|
||||
) => {
|
||||
stats.recovered += 1;
|
||||
}
|
||||
Ok(
|
||||
TransitionTransactionRecoveryOutcome::Retained
|
||||
| TransitionTransactionRecoveryOutcome::RetainedAmbiguous(_)
|
||||
| TransitionTransactionRecoveryOutcome::OperatorRequired(_),
|
||||
) => {
|
||||
Ok(TransitionTransactionRecoveryOutcome::Retained) => {
|
||||
stats.retained += 1;
|
||||
debug!(
|
||||
event = EVENT_LIFECYCLE_TRANSITION_TRANSACTION_RECOVERY,
|
||||
@@ -1982,74 +1509,11 @@ fn state_requires_known_remote_version(state: TransitionTransactionState) -> boo
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::HashMap;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
||||
use super::*;
|
||||
|
||||
const BACKEND_FINGERPRINT: [u8; 32] = [7; 32];
|
||||
|
||||
struct RecoveryAttemptDropGuard(Arc<AtomicBool>);
|
||||
|
||||
impl Drop for RecoveryAttemptDropGuard {
|
||||
fn drop(&mut self) {
|
||||
self.0.store(true, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
|
||||
async fn pending_recovery_attempt(started: Arc<tokio::sync::Notify>, dropped: Arc<AtomicBool>) -> EcstoreResult<()> {
|
||||
let _drop_guard = RecoveryAttemptDropGuard(dropped);
|
||||
started.notify_one();
|
||||
std::future::pending().await
|
||||
}
|
||||
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn transition_recovery_timeout_and_cancellation_drop_inflight_attempts() {
|
||||
let timeout_started = Arc::new(tokio::sync::Notify::new());
|
||||
let timeout_dropped = Arc::new(AtomicBool::new(false));
|
||||
let timeout_task = tokio::spawn({
|
||||
let started = Arc::clone(&timeout_started);
|
||||
let dropped = Arc::clone(&timeout_dropped);
|
||||
async move {
|
||||
await_transition_transaction_recovery(
|
||||
&CancellationToken::new(),
|
||||
TRANSITION_TRANSACTION_RECOVERY_TIMEOUT,
|
||||
pending_recovery_attempt(started, dropped),
|
||||
)
|
||||
.await
|
||||
}
|
||||
});
|
||||
timeout_started.notified().await;
|
||||
tokio::time::advance(TRANSITION_TRANSACTION_RECOVERY_TIMEOUT).await;
|
||||
let timed_out = timeout_task.await.expect("timeout wrapper task should join");
|
||||
assert!(matches!(timed_out, Some(Err(_))), "outer timeout should fail the recovery pass");
|
||||
assert!(timeout_dropped.load(Ordering::SeqCst), "outer timeout must drop its in-flight attempt");
|
||||
|
||||
let cancel_token = CancellationToken::new();
|
||||
let cancel_started = Arc::new(tokio::sync::Notify::new());
|
||||
let cancel_dropped = Arc::new(AtomicBool::new(false));
|
||||
let cancel_task = tokio::spawn({
|
||||
let cancel_token = cancel_token.clone();
|
||||
let started = Arc::clone(&cancel_started);
|
||||
let dropped = Arc::clone(&cancel_dropped);
|
||||
async move {
|
||||
await_transition_transaction_recovery(
|
||||
&cancel_token,
|
||||
TRANSITION_TRANSACTION_RECOVERY_TIMEOUT,
|
||||
pending_recovery_attempt(started, dropped),
|
||||
)
|
||||
.await
|
||||
}
|
||||
});
|
||||
cancel_started.notified().await;
|
||||
cancel_token.cancel();
|
||||
let cancelled = cancel_task.await.expect("cancellation wrapper task should join");
|
||||
assert!(cancelled.is_none(), "outer cancellation should stop the recovery loop");
|
||||
assert!(
|
||||
cancel_dropped.load(Ordering::SeqCst),
|
||||
"outer cancellation must drop its in-flight attempt"
|
||||
);
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct MemoryTransactionStore {
|
||||
records: HashMap<Uuid, Vec<u8>>,
|
||||
@@ -2504,19 +1968,5 @@ mod tests {
|
||||
transition_transaction_record_object_name(Uuid::nil()),
|
||||
Err(TransitionTransactionError::Corrupt("transaction_id is nil"))
|
||||
));
|
||||
assert_eq!(
|
||||
transition_transaction_id_from_record_object_name(&object).expect("canonical record path should parse"),
|
||||
transaction_id
|
||||
);
|
||||
for malformed in [
|
||||
object.to_ascii_uppercase(),
|
||||
object.replace("/aa/aa/", "/ff/aa/"),
|
||||
object.replace("/aa/aa/", "/aa/aa/extra/"),
|
||||
] {
|
||||
assert!(matches!(
|
||||
transition_transaction_id_from_record_object_name(&malformed),
|
||||
Err(TransitionTransactionError::Corrupt(_))
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1611,6 +1611,34 @@ mod test {
|
||||
assert!(bm.bucket_target_config.is_none());
|
||||
}
|
||||
|
||||
/// rustfs/backlog#2309: the MinIO-origin `.metadata.bin` this repository
|
||||
/// already carries as a compatibility fixture stores
|
||||
/// `BucketTargetsConfigJSON` as a bare JSON array, which `BucketTargets`
|
||||
/// (a `{"targets":[…]}` struct with no array fallback) cannot decode. The
|
||||
/// bytes below are the exact payload the fixture in
|
||||
/// `metadata_test.rs::TEST_BUCKET_METADATA_HEX` decodes to, so if RustFS
|
||||
/// ever grows the array-shaped compatibility parse, this test is where the
|
||||
/// upgrade break is pinned and where the decision has to be recorded.
|
||||
#[test]
|
||||
fn minio_array_shaped_bucket_targets_are_unreadable() {
|
||||
let minio_array = br#"[{"endpoint":"http://target.example.com","targetBucket":"tb","region":"us-east-1"}]"#.to_vec();
|
||||
let mut bm = BucketMetadata::new("minio-array-targets");
|
||||
bm.bucket_targets_config_json = minio_array.clone();
|
||||
|
||||
bm.parse_all_configs()
|
||||
.expect("a MinIO-shaped targets blob must not fail the whole metadata load");
|
||||
|
||||
assert!(
|
||||
bm.bucket_targets_unreadable(),
|
||||
"an array-shaped MinIO targets blob is unreadable, not an empty target set"
|
||||
);
|
||||
assert!(bm.bucket_target_config.is_none());
|
||||
assert_eq!(
|
||||
bm.bucket_targets_config_json, minio_array,
|
||||
"the raw MinIO bytes must survive so the configuration stays recoverable"
|
||||
);
|
||||
}
|
||||
|
||||
/// The invariant every branch of `parse_all_configs` shares: a stored but
|
||||
/// undecodable payload keeps its raw bytes and leaves the typed field
|
||||
/// `None`, so no branch fabricates a value. What a reader may then do with
|
||||
|
||||
@@ -825,11 +825,6 @@ mod tests {
|
||||
manual_transition_scope_record_object_name, manual_transition_task_object_name,
|
||||
manual_transition_worker_result_object_name, manual_transition_worker_result_task_key,
|
||||
},
|
||||
recovery_control::{
|
||||
IlmRecoveryClassification, IlmRecoveryControl, IlmRecoveryControlIdentity, IlmRecoveryErrorCode,
|
||||
IlmRecoveryProtocol, MAX_RECOVERY_ATTEMPTS, load_recovery_control, observe_recovery_source,
|
||||
save_recovery_control_if_absent,
|
||||
},
|
||||
tier_delete_journal::{
|
||||
DecommissionCheckpointTargetFailureHook, TIER_DELETE_DISPATCH_MANIFEST_PREFIX, TIER_DELETE_JOURNAL_PREFIX,
|
||||
TierDeleteChunkTestBarrier, TierDeleteChunkTestStage, TierDeleteDispatchBatchLimitGuard,
|
||||
@@ -849,13 +844,12 @@ mod tests {
|
||||
},
|
||||
transition_transaction::{
|
||||
TRANSITION_TRANSACTION_RECORD_PREFIX, TransitionCleanupDecision, TransitionCleanupProof, TransitionOperatorError,
|
||||
TransitionOperatorProbe, TransitionRecoveryClaimBarrier, TransitionRecoveryTerminalBarrier,
|
||||
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
|
||||
TransitionTransactionInit, TransitionTransactionState, delete_transition_candidate_for_operator,
|
||||
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
|
||||
load_transition_transaction_record, recover_transition_transaction_records,
|
||||
recover_transition_transaction_records_at, save_transition_transaction_record,
|
||||
save_transition_transaction_record_if_current, transition_recovery_control_id,
|
||||
TransitionOperatorProbe, TransitionRecoveryClaimBarrier, TransitionRemoteVersion, TransitionSourceIdentity,
|
||||
TransitionSourceVersionMode, TransitionTransaction, TransitionTransactionInit, TransitionTransactionState,
|
||||
delete_transition_candidate_for_operator, finalize_missing_transition_transaction_for_operator,
|
||||
inspect_transition_transaction_for_operator, load_transition_transaction_record,
|
||||
recover_transition_transaction_records, recover_transition_transaction_records_at,
|
||||
save_transition_transaction_record, save_transition_transaction_record_if_current,
|
||||
transition_transaction_record_object_name,
|
||||
},
|
||||
validate_durable_ilm_record,
|
||||
@@ -19245,105 +19239,6 @@ mod tests {
|
||||
assert!(!Arc::ptr_eq(&ctx_a, &ctx_b), "the regression requires two distinct instance contexts");
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn transition_transaction_recovery_expires_abandoned_attempt_at_budget_bound() {
|
||||
let temp_dir = tempfile::tempdir().expect("create temp store dir");
|
||||
let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store(
|
||||
temp_dir.path(),
|
||||
"transition-transaction-expired-attempt-budget",
|
||||
&[4],
|
||||
))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let transaction = TransitionTransaction::new(TransitionTransactionInit {
|
||||
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
|
||||
transaction_id: uuid::Uuid::new_v4(),
|
||||
owner_epoch: uuid::Uuid::new_v4(),
|
||||
write_id: uuid::Uuid::new_v4(),
|
||||
source: TransitionSourceIdentity {
|
||||
bucket: "source-bucket".to_string(),
|
||||
object: "source-object".to_string(),
|
||||
version_id: Some(uuid::Uuid::new_v4()),
|
||||
data_dir: uuid::Uuid::new_v4(),
|
||||
mod_time_unix_nanos: 1_770_000_000_000_000_000,
|
||||
size: 42,
|
||||
etag: "source-etag".to_string(),
|
||||
version_mode: TransitionSourceVersionMode::Versioned,
|
||||
},
|
||||
tier_name: "UNUSEDABANDONEDTIER".to_string(),
|
||||
backend_fingerprint: [7; 32],
|
||||
not_after_unix_nanos: 1,
|
||||
})
|
||||
.expect("transaction should build");
|
||||
save_transition_transaction_record(store.clone(), &transaction)
|
||||
.await
|
||||
.expect("transaction record should persist");
|
||||
|
||||
let record_name =
|
||||
transition_transaction_record_object_name(transaction.transaction_id).expect("transaction record name should derive");
|
||||
let source = observe_recovery_source(
|
||||
store.clone(),
|
||||
&record_name,
|
||||
crate::bucket::lifecycle::transition_transaction::TRANSITION_TRANSACTION_SCHEMA,
|
||||
)
|
||||
.await
|
||||
.expect("transaction source generation should be observable");
|
||||
let mut control = IlmRecoveryControl::new(
|
||||
IlmRecoveryControlIdentity {
|
||||
protocol: IlmRecoveryProtocol::TransitionTransaction,
|
||||
canonical_source_path: record_name,
|
||||
stable_operation_identity: transaction.transaction_id.to_string(),
|
||||
record_class: "transition_transaction_v1".to_string(),
|
||||
},
|
||||
source.generation,
|
||||
IlmRecoveryClassification::Retrying,
|
||||
2_000_000_000,
|
||||
IlmRecoveryErrorCode::None,
|
||||
)
|
||||
.expect("recovery control should build");
|
||||
let mut now = 3_000_000_000;
|
||||
for _ in 1..MAX_RECOVERY_ATTEMPTS {
|
||||
now = now.max(control.next_attempt_at_unix_nanos.unwrap_or(now));
|
||||
control
|
||||
.claim("cancelled-or-timed-out-owner", uuid::Uuid::new_v4(), now, 1)
|
||||
.expect("abandoned attempt should claim");
|
||||
control
|
||||
.record_expired_attempt(now + 1)
|
||||
.expect("expired attempt should consume retry budget");
|
||||
now += 2;
|
||||
}
|
||||
now = now.max(control.next_attempt_at_unix_nanos.expect("last retry should have a backoff"));
|
||||
control
|
||||
.claim("cancelled-or-timed-out-owner", uuid::Uuid::new_v4(), now, 1)
|
||||
.expect("final abandoned attempt should claim");
|
||||
save_recovery_control_if_absent(store.clone(), &control)
|
||||
.await
|
||||
.expect("claimed recovery control should persist");
|
||||
|
||||
let stats = recover_transition_transaction_records_at(store.clone(), 100, None, i128::from(now + 1))
|
||||
.await
|
||||
.expect("recovery should account for the expired attempt");
|
||||
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 0, 1, 0));
|
||||
|
||||
let control_id = transition_recovery_control_id(&transaction).expect("control id should derive");
|
||||
let persisted = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &control_id)
|
||||
.await
|
||||
.expect("expired recovery control should remain inspectable");
|
||||
assert_eq!(persisted.control.classification, IlmRecoveryClassification::OperatorRequired);
|
||||
assert_eq!(persisted.control.attempt_count, u64::from(MAX_RECOVERY_ATTEMPTS));
|
||||
assert_eq!(persisted.control.consecutive_failure_count, MAX_RECOVERY_ATTEMPTS);
|
||||
assert_eq!(persisted.control.last_error_code, IlmRecoveryErrorCode::AttemptLeaseExpired);
|
||||
assert!(persisted.control.owner.is_none());
|
||||
assert_eq!(
|
||||
transition_transaction_record_count(store).await,
|
||||
1,
|
||||
"budget exhaustion must retain the source record"
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
@@ -19456,7 +19351,6 @@ mod tests {
|
||||
),
|
||||
];
|
||||
let mut expected_removes = Vec::new();
|
||||
let mut recovery_control_ids = Vec::new();
|
||||
for (case, put_version, remote_version, source_mode) in cases {
|
||||
let mut transaction = TransitionTransaction::new(TransitionTransactionInit {
|
||||
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
|
||||
@@ -19494,8 +19388,6 @@ mod tests {
|
||||
save_transition_transaction_record(store.clone(), &transaction)
|
||||
.await
|
||||
.expect("transaction record should persist");
|
||||
recovery_control_ids
|
||||
.push(transition_recovery_control_id(&transaction).expect("transition recovery control id should derive"));
|
||||
expected_removes.push((transaction.remote_object, put_version));
|
||||
}
|
||||
|
||||
@@ -19511,14 +19403,6 @@ mod tests {
|
||||
assert_eq!(actual_removes, expected_removes, "recovery must preserve each remote version shape");
|
||||
assert_eq!(backend.exact_remove_count(), 2);
|
||||
assert_eq!(backend.object_count().await, 0);
|
||||
for control_id in recovery_control_ids {
|
||||
let control = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &control_id)
|
||||
.await
|
||||
.expect("completed recovery control should remain inspectable");
|
||||
assert_eq!(control.control.classification, IlmRecoveryClassification::Terminal);
|
||||
assert_eq!(control.control.attempt_count, 1);
|
||||
assert!(control.control.owner.is_none());
|
||||
}
|
||||
|
||||
let replay = recover_transition_transaction_records(store, 100, None)
|
||||
.await
|
||||
@@ -19531,94 +19415,6 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn transition_transaction_recovery_resumes_source_cleanup_after_terminal_crash() {
|
||||
let temp_dir = tempfile::tempdir().expect("create temp store dir");
|
||||
let (ctx, store, _shutdown) =
|
||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "transition-transaction-terminal-crash", &[4]))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let tier_name = "TXTERMINALCRASH";
|
||||
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
|
||||
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
|
||||
.await
|
||||
.expect("tier lease should resolve")
|
||||
.backend_identity();
|
||||
let remote_version = uuid::Uuid::new_v4().to_string();
|
||||
let mut transaction = TransitionTransaction::new(TransitionTransactionInit {
|
||||
deployment_id: ctx.deployment_id().expect("test store should initialize deployment id"),
|
||||
transaction_id: uuid::Uuid::new_v4(),
|
||||
owner_epoch: uuid::Uuid::new_v4(),
|
||||
write_id: uuid::Uuid::new_v4(),
|
||||
source: TransitionSourceIdentity {
|
||||
bucket: "source-bucket".to_string(),
|
||||
object: "source-object".to_string(),
|
||||
version_id: Some(uuid::Uuid::new_v4()),
|
||||
data_dir: uuid::Uuid::new_v4(),
|
||||
mod_time_unix_nanos: 1_770_000_000_000_000_000,
|
||||
size: 42,
|
||||
etag: "source-etag".to_string(),
|
||||
version_mode: TransitionSourceVersionMode::Versioned,
|
||||
},
|
||||
tier_name: tier_name.to_string(),
|
||||
backend_fingerprint: backend_identity,
|
||||
not_after_unix_nanos: 1,
|
||||
})
|
||||
.expect("transaction should build");
|
||||
transaction
|
||||
.advance(
|
||||
transaction.fence(),
|
||||
TransitionTransactionState::Uploaded,
|
||||
Some(TransitionRemoteVersion::versioned(remote_version.clone())),
|
||||
)
|
||||
.expect("transaction should enter uploaded state");
|
||||
backend.set_put_remote_version(Some(remote_version)).await;
|
||||
let candidate = bytes::Bytes::from_static(b"terminal crash candidate");
|
||||
backend
|
||||
.put(
|
||||
&transaction.remote_object,
|
||||
ReaderImpl::Body(candidate.clone()),
|
||||
i64::try_from(candidate.len()).expect("test candidate length should fit i64"),
|
||||
)
|
||||
.await
|
||||
.expect("mock backend should accept candidate");
|
||||
save_transition_transaction_record(store.clone(), &transaction)
|
||||
.await
|
||||
.expect("transaction record should persist");
|
||||
let control_id = transition_recovery_control_id(&transaction).expect("control id should derive");
|
||||
|
||||
let barrier = TransitionRecoveryTerminalBarrier::install(transaction.transaction_id);
|
||||
let recovery_store = store.clone();
|
||||
let recovery = tokio::spawn(async move { recover_transition_transaction_records(recovery_store, 100, None).await });
|
||||
barrier.wait_until_paused().await;
|
||||
let terminal = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &control_id)
|
||||
.await
|
||||
.expect("terminal control should persist before source cleanup");
|
||||
assert_eq!(terminal.control.classification, IlmRecoveryClassification::Terminal);
|
||||
assert_eq!(transition_transaction_record_count(store.clone()).await, 1);
|
||||
assert_eq!(backend.object_count().await, 0);
|
||||
assert_eq!(backend.exact_remove_count(), 1);
|
||||
|
||||
recovery.abort();
|
||||
assert!(
|
||||
recovery
|
||||
.await
|
||||
.expect_err("recovery should be cancelled at the crash boundary")
|
||||
.is_cancelled()
|
||||
);
|
||||
drop(barrier);
|
||||
|
||||
let replay = recover_transition_transaction_records(store.clone(), 100, None)
|
||||
.await
|
||||
.expect("terminal control should resume source cleanup without another remote delete");
|
||||
assert_eq!((replay.scanned, replay.recovered, replay.retained, replay.failed), (1, 1, 0, 0));
|
||||
assert_eq!(transition_transaction_record_count(store).await, 0);
|
||||
assert_eq!(backend.exact_remove_count(), 1, "terminal replay must not repeat the remote delete");
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
@@ -19676,8 +19472,6 @@ mod tests {
|
||||
save_transition_transaction_record(store.clone(), &uploaded)
|
||||
.await
|
||||
.expect("transaction record should persist");
|
||||
let recovery_control_id =
|
||||
transition_recovery_control_id(&uploaded).expect("transition recovery control id should derive");
|
||||
|
||||
let barrier = TransitionRecoveryClaimBarrier::install(uploaded.transaction_id);
|
||||
let recovery_store = store.clone();
|
||||
@@ -19699,17 +19493,11 @@ mod tests {
|
||||
.expect("recovery should treat the lost CAS as a retained transaction");
|
||||
assert_eq!((stats.scanned, stats.recovered, stats.retained, stats.failed), (1, 0, 1, 0));
|
||||
assert_eq!(
|
||||
load_transition_transaction_record(store.clone(), uploaded.transaction_id)
|
||||
load_transition_transaction_record(store, uploaded.transaction_id)
|
||||
.await
|
||||
.expect("newer transaction revision must remain"),
|
||||
active
|
||||
);
|
||||
let control = load_recovery_control(store, IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
|
||||
.await
|
||||
.expect("lost source CAS should retain a retryable recovery control");
|
||||
assert_eq!(control.control.classification, IlmRecoveryClassification::Retrying);
|
||||
assert_eq!(control.control.consecutive_failure_count, 1);
|
||||
assert_eq!(control.control.last_error_code, IlmRecoveryErrorCode::SourceGenerationChanged);
|
||||
assert_eq!(backend.object_count().await, 1, "a stale recovery must not delete the candidate");
|
||||
assert_eq!(backend.remove_count().await, 0);
|
||||
}
|
||||
@@ -20011,15 +19799,27 @@ mod tests {
|
||||
not_after_unix_nanos: 1_780_000_000_000_000_000,
|
||||
})
|
||||
.expect("transaction should build");
|
||||
transaction
|
||||
let uploaded_fence = transaction
|
||||
.advance(
|
||||
transaction.fence(),
|
||||
TransitionTransactionState::Uploaded,
|
||||
Some(TransitionRemoteVersion::versioned(remote_version.clone())),
|
||||
Some(TransitionRemoteVersion::versioned(remote_version)),
|
||||
)
|
||||
.expect("transaction should enter uploaded state");
|
||||
transaction
|
||||
.mark_cleanup_pending(
|
||||
uploaded_fence,
|
||||
TransitionCleanupProof {
|
||||
transaction_id: transaction.transaction_id,
|
||||
write_id: transaction.write_id,
|
||||
remote_object: transaction.remote_object.clone(),
|
||||
remote_version: transaction.remote_version.clone(),
|
||||
backend_fingerprint: transaction.backend_fingerprint,
|
||||
decision: TransitionCleanupDecision::UploadAbortedBeforeLocalCommit,
|
||||
},
|
||||
)
|
||||
.expect("transaction should enter cleanup pending state");
|
||||
let candidate = bytes::Bytes::from_static(b"cleanup pending candidate retained after failure");
|
||||
backend.set_put_remote_version(Some(remote_version)).await;
|
||||
backend
|
||||
.put(
|
||||
&transaction.remote_object,
|
||||
@@ -20031,8 +19831,6 @@ mod tests {
|
||||
save_transition_transaction_record(store.clone(), &transaction)
|
||||
.await
|
||||
.expect("transaction record should persist");
|
||||
let recovery_control_id =
|
||||
transition_recovery_control_id(&transaction).expect("transition recovery control id should derive");
|
||||
|
||||
backend.set_remove_failure(true);
|
||||
let stats = recover_transition_transaction_records(store.clone(), 100, None)
|
||||
@@ -20048,42 +19846,6 @@ mod tests {
|
||||
assert_eq!(backend.remove_versions().await, Vec::<(String, String)>::new());
|
||||
assert_eq!(backend.exact_remove_count(), 1);
|
||||
assert_eq!(backend.object_count().await, 1);
|
||||
let control = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
|
||||
.await
|
||||
.expect("failed recovery control should persist");
|
||||
assert_eq!(control.control.classification, IlmRecoveryClassification::Retrying);
|
||||
assert_eq!(control.control.attempt_count, 1);
|
||||
assert_eq!(control.control.consecutive_failure_count, 1);
|
||||
assert!(
|
||||
control
|
||||
.control
|
||||
.next_attempt_at_unix_nanos
|
||||
.is_some_and(|next| next > OffsetDateTime::now_utc().unix_timestamp_nanos() as i64)
|
||||
);
|
||||
|
||||
backend.set_remove_failure(false);
|
||||
let replay = recover_transition_transaction_records(store.clone(), 100, None)
|
||||
.await
|
||||
.expect("recovery before the persisted deadline should be skipped");
|
||||
assert_eq!((replay.scanned, replay.recovered, replay.retained, replay.failed), (1, 0, 1, 0));
|
||||
assert_eq!(backend.exact_remove_count(), 1, "persisted backoff must prevent an immediate retry");
|
||||
|
||||
let retry_at = control
|
||||
.control
|
||||
.next_attempt_at_unix_nanos
|
||||
.expect("retry deadline should persist");
|
||||
let retried = recover_transition_transaction_records_at(store.clone(), 100, None, i128::from(retry_at) + 1)
|
||||
.await
|
||||
.expect("recovery at the persisted deadline should retry the advanced source generation");
|
||||
assert_eq!((retried.scanned, retried.recovered, retried.retained, retried.failed), (1, 1, 0, 0));
|
||||
assert_eq!(backend.exact_remove_count(), 2);
|
||||
assert_eq!(backend.object_count().await, 0);
|
||||
assert_eq!(transition_transaction_record_count(store.clone()).await, 0);
|
||||
let terminal = load_recovery_control(store, IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
|
||||
.await
|
||||
.expect("completed retry control should remain inspectable");
|
||||
assert_eq!(terminal.control.classification, IlmRecoveryClassification::Terminal);
|
||||
assert_eq!(terminal.control.attempt_count, 2);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
@@ -20384,10 +20146,6 @@ mod tests {
|
||||
local_commit_started
|
||||
.advance(local_commit_started.fence(), TransitionTransactionState::LocalCommitStarted, None)
|
||||
.expect("transaction should enter local commit state");
|
||||
let upload_started_control_id =
|
||||
transition_recovery_control_id(&upload_started).expect("upload-started control id should derive");
|
||||
let local_commit_control_id =
|
||||
transition_recovery_control_id(&local_commit_started).expect("local-commit control id should derive");
|
||||
|
||||
backend.set_put_remote_version(Some(remote_version)).await;
|
||||
for transaction in [&upload_started, &local_commit_started] {
|
||||
@@ -20418,19 +20176,6 @@ mod tests {
|
||||
assert_eq!(backend.object_count().await, 2, "recovery must not delete an unproven remote candidate");
|
||||
assert_eq!(backend.remove_count().await, 0);
|
||||
assert_eq!(backend.exact_remove_count(), 0);
|
||||
let upload_started_control =
|
||||
load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &upload_started_control_id)
|
||||
.await
|
||||
.expect("upload-started control should persist");
|
||||
assert_eq!(
|
||||
upload_started_control.control.classification,
|
||||
IlmRecoveryClassification::RetainedAmbiguous
|
||||
);
|
||||
let local_commit_control =
|
||||
load_recovery_control(store, IlmRecoveryProtocol::TransitionTransaction, &local_commit_control_id)
|
||||
.await
|
||||
.expect("local-commit control should persist");
|
||||
assert_eq!(local_commit_control.control.classification, IlmRecoveryClassification::OperatorRequired);
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
@@ -20781,12 +20526,6 @@ mod tests {
|
||||
"an unsupported provider probe must retain the unknown upload"
|
||||
);
|
||||
assert_eq!(transition_transaction_record_count(store.clone()).await, 1);
|
||||
let recovery_control_id =
|
||||
transition_recovery_control_id(&transaction).expect("transition recovery control id should derive");
|
||||
let control = load_recovery_control(store.clone(), IlmRecoveryProtocol::TransitionTransaction, &recovery_control_id)
|
||||
.await
|
||||
.expect("unsupported probe control should persist");
|
||||
assert_eq!(control.control.classification, IlmRecoveryClassification::RetainedAmbiguous);
|
||||
assert!(
|
||||
backend.contains(&transaction.remote_object).await,
|
||||
"unsupported recovery must not delete the candidate"
|
||||
|
||||
@@ -1020,7 +1020,6 @@ mod serial_tests {
|
||||
}
|
||||
|
||||
let (_disk_paths, ecstore) = setup_isolated_test_env(false).await;
|
||||
let expired_recovery_time = i128::from(i64::MAX / 2);
|
||||
|
||||
for case in [
|
||||
CleanupCase::Persisted,
|
||||
@@ -1118,7 +1117,7 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("active unknown ownership must remain fenced after the transaction store was offline");
|
||||
assert_eq!((retained.scanned, retained.recovered, retained.retained, retained.failed), (1, 0, 1, 0));
|
||||
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, expired_recovery_time)
|
||||
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, i128::MAX)
|
||||
.await
|
||||
.expect("expired unknown ownership may use the provider's missing proof");
|
||||
assert_eq!(
|
||||
@@ -1167,7 +1166,7 @@ mod serial_tests {
|
||||
assert_eq!(retained.recovered, 0);
|
||||
assert_eq!(retained.retained + retained.failed, 1);
|
||||
backend.set_remove_failure(false);
|
||||
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, expired_recovery_time)
|
||||
let recovered = recover_transition_transaction_records_at(ecstore.clone(), 100, None, i128::MAX)
|
||||
.await
|
||||
.expect("expired recovery should delete the candidate after the backend becomes available");
|
||||
assert_eq!(
|
||||
|
||||
@@ -1,12 +1,12 @@
|
||||
# Bucket metadata diagnostics and recovery
|
||||
|
||||
`GET /rustfs/admin/v3/export-bucket-metadata` keeps its strict behavior: an unreadable configuration fails the export. The optional `bucket` query selects one bucket; omitting it selects all buckets.
|
||||
`GET /rustfs/admin/v3/export-bucket-metadata` exports every configuration it can read. A configuration that is stored but unreadable is never exported and never replaced by a fabricated default; instead the bucket gains an entry `<bucket>/rustfs-unreadable-configs.json` of the shape `{"bucket": …, "unreadable": [{"config": …, "error": …}]}` naming each configuration that could not be read and why, and the export continues, so one bucket's undecodable payload cannot cost an operator the whole-cluster backup (rustfs/backlog#2309). The marker name is outside the configuration namespace importers dispatch on, so importing the archive back leaves the affected bucket's stored bytes untouched. A failure of the server's own serialization or archive writing still fails the export. The optional `bucket` query selects one bucket; omitting it selects all buckets.
|
||||
|
||||
To inspect readable configurations while identifying failures, use the same authenticated endpoint with `?diagnostic=true`. This requires the existing `ExportBucketMetadataAction` permission. A successful response has:
|
||||
To collect a shareable support artifact that identifies the failures without carrying parser detail, use the same authenticated endpoint with `?diagnostic=true`. This requires the existing `ExportBucketMetadataAction` permission. A successful response has:
|
||||
|
||||
- Filename `bucket-meta-diagnostic.zip` and header `x-rustfs-bucket-metadata-export: diagnostic`.
|
||||
- Readable entries under `_diagnostic/<bucket>/<config>`; target credentials remain redacted.
|
||||
- `_diagnostic-manifest.json`, containing `version: 1`, `mode: "diagnostic"`, `complete`, and an `errors` array. Each error identifies `bucket`, `config`, and the fixed code `configuration_unavailable`. The archive excludes unreadable payloads and parser error details.
|
||||
- `_diagnostic-manifest.json`, containing `version: 1`, `mode: "diagnostic"`, `complete`, and an `errors` array. Each error identifies `bucket`, `config`, and the fixed code `configuration_unavailable`. The archive excludes unreadable payloads and parser error details, and therefore carries no `rustfs-unreadable-configs.json` marker: the manifest reports the same failures with less detail, which is what makes a diagnostic archive safe to hand out.
|
||||
|
||||
`complete` reports whether all supported configuration reads succeeded. A diagnostic archive is never a restorable backup, including when `complete` is true. Import rejects the manifest or reserved directory before any bucket creation or configuration write. The reserved directory is not a valid bucket name, so older importers cannot restore diagnostic entries as ordinary bucket configurations.
|
||||
|
||||
@@ -14,9 +14,11 @@ To inspect readable configurations while identifying failures, use the same auth
|
||||
|
||||
RustFS currently accepts the documented `{"targets": [...]}` object format. It cannot decrypt MinIO KMS-encrypted target metadata. Unreadable target payloads remain failures instead of being interpreted as an empty target set; diagnostic export and replacement import do not add MinIO KMS decryption support.
|
||||
|
||||
1. Inspect the diagnostic manifest to identify affected buckets. Preserve a separate backup of the original source configuration and any credentials needed for recovery.
|
||||
1. Inspect the diagnostic manifest, or the `rustfs-unreadable-configs.json` marker in an ordinary export, to identify affected buckets. Preserve a separate backup of the original source configuration and any credentials needed for recovery.
|
||||
2. Prepare a ZIP containing `<bucket>/bucket-targets.json` with a valid RustFS replacement, whose top-level shape is `{"targets": [...]}`. Supply the intended target settings and credentials; exported credentials are redacted. Use `{"targets": []}` only when intentionally clearing all targets, and reconcile any replication rules that reference removed targets.
|
||||
3. Submit the ZIP to the existing authenticated `PUT /rustfs/admin/v3/import-bucket-metadata` endpoint with `ImportBucketMetadataAction` permission. Import validates the replacement and persists it against the bucket incarnation; it does not need to parse the old target payload successfully.
|
||||
4. Verify target listing and the intended replication configuration. Retry the ordinary strict metadata export to confirm the unreadable configuration no longer blocks it.
|
||||
4. Verify target listing and the intended replication configuration. Retry the ordinary metadata export to confirm the bucket no longer carries an unreadable marker.
|
||||
|
||||
Alternatively, `PUT /rustfs/admin/v3/set-remote-target?replace-unreadable=true` discards an undecodable target set as part of setting a replacement target. The flag is the operator's explicit acknowledgement that the stored set is being thrown away; without it the request is refused rather than rewriting an unreadable set from a partial view.
|
||||
|
||||
Do not submit the diagnostic archive itself to the import endpoint. Copy only reviewed replacement entries into an ordinary import archive.
|
||||
|
||||
@@ -67,6 +67,17 @@ use zip::{ZipArchive, ZipWriter, write::SimpleFileOptions};
|
||||
const DIAGNOSTIC_EXPORT_PREFIX: &str = "_diagnostic";
|
||||
const DIAGNOSTIC_EXPORT_MANIFEST: &str = "_diagnostic-manifest.json";
|
||||
|
||||
/// Archive entry naming the configurations a bucket stores but this build
|
||||
/// could not read (rustfs/backlog#2309).
|
||||
///
|
||||
/// The name deliberately sits outside the configuration-file namespace the
|
||||
/// importers switch on — `ImportBucketMetadata` matches known configuration
|
||||
/// names and ignores everything else — so no importer can mistake the marker
|
||||
/// for a configuration. Ordinary exports carry it; diagnostic exports report
|
||||
/// the same failures through [`DIAGNOSTIC_EXPORT_MANIFEST`] instead, which
|
||||
/// deliberately withholds the parser detail this marker records.
|
||||
const EXPORT_UNREADABLE_MANIFEST: &str = "rustfs-unreadable-configs.json";
|
||||
|
||||
const LOG_COMPONENT_ADMIN: &str = "admin";
|
||||
const LOG_SUBSYSTEM_BUCKET_META: &str = "bucket_meta";
|
||||
const EVENT_ADMIN_BUCKET_META_STATE: &str = "admin_bucket_meta_state";
|
||||
@@ -76,31 +87,82 @@ fn export_internal_error(message: impl Into<String>) -> s3s::S3Error {
|
||||
s3_error!(InternalError, "{message}")
|
||||
}
|
||||
|
||||
fn checked_raw_xml<T, E, F>(validated: &T, raw: Vec<u8>, parse: F) -> S3Result<Vec<u8>>
|
||||
/// One configuration that is stored for a bucket but could not be exported.
|
||||
#[derive(serde::Serialize)]
|
||||
struct UnreadableExportEntry {
|
||||
config: &'static str,
|
||||
error: String,
|
||||
}
|
||||
|
||||
#[derive(serde::Serialize)]
|
||||
struct UnreadableExportManifest<'a> {
|
||||
bucket: &'a str,
|
||||
unreadable: &'a [UnreadableExportEntry],
|
||||
}
|
||||
|
||||
/// Why one of a bucket's configurations could not be exported.
|
||||
///
|
||||
/// The two variants are what an ordinary export dispatches on: a bucket whose
|
||||
/// stored bytes this build cannot turn into a configuration is named and
|
||||
/// skipped, while a failure of our own output machinery still fails the whole
|
||||
/// export closed.
|
||||
#[derive(Debug)]
|
||||
enum ExportConfigError {
|
||||
/// The configuration is stored but this build cannot read it: a
|
||||
/// MinIO-origin or otherwise undecodable blob, or a revision that moved
|
||||
/// underneath the export. No retry of ours turns those bytes into a
|
||||
/// configuration, so one such bucket must not abort a whole-cluster export
|
||||
/// (rustfs/backlog#2309).
|
||||
Unreadable(String),
|
||||
/// This build failed to produce its own output for a configuration it had
|
||||
/// already decoded. Nothing about the stored bytes is in doubt, so the
|
||||
/// export fails closed rather than reporting healthy metadata as
|
||||
/// unreadable.
|
||||
Internal(s3s::S3Error),
|
||||
}
|
||||
|
||||
impl ExportConfigError {
|
||||
fn unreadable(message: impl Into<String>) -> Self {
|
||||
Self::Unreadable(message.into())
|
||||
}
|
||||
|
||||
fn internal(message: impl Into<String>) -> Self {
|
||||
Self::Internal(export_internal_error(message))
|
||||
}
|
||||
}
|
||||
|
||||
fn checked_raw_xml<T, E, F>(validated: &T, raw: Vec<u8>, parse: F) -> Result<Vec<u8>, ExportConfigError>
|
||||
where
|
||||
T: PartialEq,
|
||||
E: std::fmt::Display,
|
||||
F: FnOnce(&[u8]) -> Result<T, E>,
|
||||
{
|
||||
let selected = parse(&raw)
|
||||
.map_err(|e| export_internal_error(format!("persisted bucket metadata changed to invalid XML during export: {e}")))?;
|
||||
let selected = parse(&raw).map_err(|e| {
|
||||
ExportConfigError::unreadable(format!("persisted bucket metadata changed to invalid XML during export: {e}"))
|
||||
})?;
|
||||
if selected != *validated {
|
||||
return Err(export_internal_error("bucket metadata changed during export"));
|
||||
return Err(ExportConfigError::unreadable("bucket metadata changed during export"));
|
||||
}
|
||||
Ok(raw)
|
||||
}
|
||||
|
||||
fn checked_versioning_xml(validated: &VersioningConfiguration, raw: Vec<u8>) -> S3Result<Vec<u8>> {
|
||||
fn checked_versioning_xml(validated: &VersioningConfiguration, raw: Vec<u8>) -> Result<Vec<u8>, ExportConfigError> {
|
||||
if raw.is_empty() {
|
||||
if *validated != VersioningConfiguration::default() {
|
||||
return Err(export_internal_error("bucket metadata changed during export"));
|
||||
return Err(ExportConfigError::unreadable("bucket metadata changed during export"));
|
||||
}
|
||||
return serialize(validated).map_err(|e| export_internal_error(format!("serialize config failed: {e}")));
|
||||
return serialize(validated).map_err(|e| ExportConfigError::internal(format!("serialize config failed: {e}")));
|
||||
}
|
||||
checked_raw_xml(validated, raw, deserialize::<VersioningConfiguration>)
|
||||
}
|
||||
|
||||
async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec<u8>>> {
|
||||
/// Bytes to export for one of a bucket's configurations.
|
||||
///
|
||||
/// `Ok(None)` means the bucket has not configured it. An `Err` never becomes an
|
||||
/// exported configuration — a fabricated default here would reach an importer
|
||||
/// as a real one — and its variant tells the caller whether the failure belongs
|
||||
/// to the stored bytes or to this build's own output; see [`ExportConfigError`].
|
||||
async fn exported_bucket_config(bucket: &str, conf: &str) -> Result<Option<Vec<u8>>, ExportConfigError> {
|
||||
match conf {
|
||||
BUCKET_POLICY_CONFIG => {
|
||||
let config: BucketPolicy = match metadata_sys::get_bucket_policy(bucket).await {
|
||||
@@ -109,11 +171,11 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "failed to load bucket metadata: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("failed to load bucket metadata: {e}")));
|
||||
}
|
||||
};
|
||||
let config_json =
|
||||
serde_json::to_vec(&config).map_err(|e| s3_error!(InternalError, "failed to serialize config: {e}"))?;
|
||||
let config_json = serde_json::to_vec(&config)
|
||||
.map_err(|e| ExportConfigError::internal(format!("failed to serialize config: {e}")))?;
|
||||
Ok(Some(config_json))
|
||||
}
|
||||
BUCKET_NOTIFICATION_CONFIG => {
|
||||
@@ -123,14 +185,14 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
Ok(None) => return Ok(None),
|
||||
};
|
||||
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("get bucket metadata failed: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")))?
|
||||
.notification_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_raw_xml(&config, raw_config, deserialize::<s3s::dto::NotificationConfiguration>)?;
|
||||
@@ -144,12 +206,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "failed to load bucket metadata: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("failed to load bucket metadata: {e}")));
|
||||
}
|
||||
};
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("failed to load bucket metadata: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("failed to load bucket metadata: {e}")))?
|
||||
.lifecycle_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_raw_xml(&config, raw_config, deserialize::<BucketLifecycleConfiguration>)?;
|
||||
@@ -163,12 +225,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "failed to load bucket metadata: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("failed to load bucket metadata: {e}")));
|
||||
}
|
||||
};
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("failed to load bucket metadata: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("failed to load bucket metadata: {e}")))?
|
||||
.tagging_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_raw_xml(&config, raw_config, deserialize::<Tagging>)?;
|
||||
@@ -182,11 +244,11 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
};
|
||||
let config_json =
|
||||
serde_json::to_vec(&config).map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?;
|
||||
serde_json::to_vec(&config).map_err(|e| ExportConfigError::internal(format!("serialize config failed: {e}")))?;
|
||||
|
||||
Ok(Some(config_json))
|
||||
}
|
||||
@@ -197,12 +259,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
};
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("get bucket metadata failed: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")))?
|
||||
.object_lock_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_raw_xml(&config, raw_config, deserialize::<ObjectLockConfiguration>)?;
|
||||
@@ -216,12 +278,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
};
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("get bucket metadata failed: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")))?
|
||||
.encryption_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_raw_xml(&config, raw_config, deserialize::<ServerSideEncryptionConfiguration>)?;
|
||||
@@ -235,12 +297,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
};
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("get bucket metadata failed: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")))?
|
||||
.versioning_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_versioning_xml(&config, raw_config)?;
|
||||
@@ -254,12 +316,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
};
|
||||
let raw_config = metadata_sys::get(bucket)
|
||||
.await
|
||||
.map_err(|e| export_internal_error(format!("get bucket metadata failed: {e}")))?
|
||||
.map_err(|e| ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")))?
|
||||
.replication_config_xml
|
||||
.clone();
|
||||
let config_xml = checked_raw_xml(&config, raw_config, deserialize::<ReplicationConfiguration>)?;
|
||||
@@ -273,12 +335,12 @@ async fn exported_bucket_config(bucket: &str, conf: &str) -> S3Result<Option<Vec
|
||||
if e == StorageError::ConfigNotFound {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(s3_error!(InternalError, "get bucket metadata failed: {e}"));
|
||||
return Err(ExportConfigError::unreadable(format!("get bucket metadata failed: {e}")));
|
||||
}
|
||||
};
|
||||
|
||||
let config_json = serde_json::to_vec(&config.redacted_credentials())
|
||||
.map_err(|e| s3_error!(InternalError, "serialize config failed: {e}"))?;
|
||||
.map_err(|e| ExportConfigError::internal(format!("serialize config failed: {e}")))?;
|
||||
|
||||
Ok(Some(config_json))
|
||||
}
|
||||
@@ -377,19 +439,52 @@ impl Operation for ExportBucketMetadata {
|
||||
];
|
||||
|
||||
for bucket in buckets {
|
||||
let mut unreadable: Vec<UnreadableExportEntry> = Vec::new();
|
||||
for &conf in confs.iter() {
|
||||
let conf_path = path_join_buf(&[bucket.name.as_str(), conf]);
|
||||
let config = match exported_bucket_config(&bucket.name, conf).await {
|
||||
Ok(Some(config)) => config,
|
||||
Ok(None) => continue,
|
||||
Err(error) if !query.diagnostic => return Err(error),
|
||||
Err(_) => {
|
||||
errors.push(serde_json::json!({
|
||||
"bucket": bucket.name,
|
||||
"config": conf,
|
||||
"code": "configuration_unavailable",
|
||||
}));
|
||||
continue;
|
||||
Err(error) => {
|
||||
if query.diagnostic {
|
||||
// A diagnostic archive names every configuration it
|
||||
// could not export under one fixed code, carrying
|
||||
// neither the payload nor the parser detail, so it
|
||||
// stays shareable (rustfs/rustfs#7225).
|
||||
errors.push(serde_json::json!({
|
||||
"bucket": bucket.name,
|
||||
"config": conf,
|
||||
"code": "configuration_unavailable",
|
||||
}));
|
||||
continue;
|
||||
}
|
||||
match error {
|
||||
// One bucket's undecodable blob must not abort the
|
||||
// whole-cluster export: record which configuration
|
||||
// could not be read and keep going, so an operator
|
||||
// migrating away still gets every readable
|
||||
// configuration (rustfs/backlog#2309).
|
||||
ExportConfigError::Unreadable(error) => {
|
||||
warn!(
|
||||
event = EVENT_ADMIN_BUCKET_META_STATE,
|
||||
component = LOG_COMPONENT_ADMIN,
|
||||
subsystem = LOG_SUBSYSTEM_BUCKET_META,
|
||||
action = "export_bucket_metadata",
|
||||
result = "config_unreadable",
|
||||
bucket = %bucket.name,
|
||||
config_name = %conf,
|
||||
error = %error,
|
||||
"admin bucket meta state"
|
||||
);
|
||||
unreadable.push(UnreadableExportEntry { config: conf, error });
|
||||
continue;
|
||||
}
|
||||
// Our own encoder failed on a configuration this
|
||||
// build had already decoded. The stored bytes are
|
||||
// not in question, so fail the export instead of
|
||||
// reporting readable metadata as unreadable.
|
||||
ExportConfigError::Internal(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
};
|
||||
let conf_path = if query.diagnostic {
|
||||
@@ -404,6 +499,23 @@ impl Operation for ExportBucketMetadata {
|
||||
.write_all(&config)
|
||||
.map_err(|e| s3_error!(InternalError, "failed to write archive entry: {e}"))?;
|
||||
}
|
||||
|
||||
// Only reachable outside diagnostic mode, which reports the same
|
||||
// failures through the archive-wide manifest instead.
|
||||
if !unreadable.is_empty() {
|
||||
let manifest = serde_json::to_vec(&UnreadableExportManifest {
|
||||
bucket: bucket.name.as_str(),
|
||||
unreadable: &unreadable,
|
||||
})
|
||||
.map_err(|e| export_internal_error(format!("failed to serialize unreadable manifest: {e}")))?;
|
||||
let manifest_path = path_join_buf(&[bucket.name.as_str(), EXPORT_UNREADABLE_MANIFEST]);
|
||||
zip_writer
|
||||
.start_file(manifest_path, SimpleFileOptions::default())
|
||||
.map_err(|e| s3_error!(InternalError, "failed to start archive entry: {e}"))?;
|
||||
zip_writer
|
||||
.write_all(&manifest)
|
||||
.map_err(|e| s3_error!(InternalError, "failed to write archive entry: {e}"))?;
|
||||
}
|
||||
}
|
||||
|
||||
if query.diagnostic {
|
||||
@@ -1364,6 +1476,10 @@ mod backup_zip_compatibility_tests {
|
||||
const ROOT_ACCESS_KEY: &str = "BUCKETMETABACKUPROOT";
|
||||
const ROOT_SECRET_KEY: &str = "bucketMetaBackupRootSecret123";
|
||||
const BUCKET: &str = "backup-compatibility";
|
||||
const UNREADABLE_BUCKET: &str = "minio-origin-targets";
|
||||
/// The exact `BucketTargetsConfigJSON` payload carried by the MinIO
|
||||
/// `.metadata.bin` fixture in `crates/ecstore/src/bucket/metadata_test.rs`.
|
||||
const MINIO_ARRAY_TARGETS: &[u8] = br#"[{"endpoint":"http://target.example.com","targetBucket":"tb","region":"us-east-1"}]"#;
|
||||
const NOTIFICATION_XML: &[u8] = b"<NotificationConfiguration>\n</NotificationConfiguration>";
|
||||
const LIFECYCLE_XML: &[u8] = b"<LifecycleConfiguration>\n<Rule><ID>expire</ID><Status>Enabled</Status><Filter><Prefix>logs/</Prefix></Filter><Expiration><Days>30</Days></Expiration></Rule>\n</LifecycleConfiguration>";
|
||||
const SSE_XML: &[u8] = b"<ServerSideEncryptionConfiguration>\n<Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule>\n</ServerSideEncryptionConfiguration>";
|
||||
@@ -1495,14 +1611,63 @@ mod backup_zip_compatibility_tests {
|
||||
.expect("publish unreadable targets fixture");
|
||||
assert!(metadata_sys::get_bucket_targets_config(UNREADABLE).await.is_err());
|
||||
|
||||
let strict_error = ExportBucketMetadata {}
|
||||
// rustfs/backlog#2309: an ordinary export no longer fails closed on
|
||||
// a configuration that is stored but unreadable. It names that one
|
||||
// configuration in the bucket's own marker entry and keeps every
|
||||
// readable configuration of every bucket, so one MinIO-origin blob
|
||||
// cannot cost an operator the whole-cluster backup.
|
||||
let ordinary = ExportBucketMetadata {}
|
||||
.call(
|
||||
admin_request(Method::GET, Uri::from_static("/rustfs/admin/v3/export-bucket-metadata"), Vec::new()),
|
||||
Params::new(),
|
||||
)
|
||||
.await
|
||||
.expect_err("a complete export must fail closed on unreadable targets");
|
||||
assert_eq!(*strict_error.code(), s3s::S3ErrorCode::InternalError);
|
||||
.expect("one unreadable configuration must not abort the ordinary export");
|
||||
assert_eq!(ordinary.output.0, StatusCode::OK);
|
||||
assert!(!ordinary.headers.contains_key("x-rustfs-bucket-metadata-export"));
|
||||
let ordinary_bytes = ordinary.output.1.collect().await.expect("read ordinary archive").to_bytes();
|
||||
let mut ordinary_archive = ZipArchive::new(Cursor::new(&ordinary_bytes)).expect("open ordinary archive");
|
||||
assert!(
|
||||
ordinary_archive
|
||||
.by_name(&format!("{HEALTHY}/{BUCKET_VERSIONING_CONFIG}"))
|
||||
.is_ok(),
|
||||
"a healthy bucket must still export while another bucket is unreadable"
|
||||
);
|
||||
assert!(
|
||||
ordinary_archive
|
||||
.by_name(&format!("{HEALTHY}/{EXPORT_UNREADABLE_MANIFEST}"))
|
||||
.is_err(),
|
||||
"a bucket whose configurations all read must carry no unreadable marker"
|
||||
);
|
||||
assert!(
|
||||
ordinary_archive
|
||||
.by_name(&format!("{UNREADABLE}/{BUCKET_TARGETS_FILE}"))
|
||||
.is_err(),
|
||||
"an unreadable targets blob must never be exported as a configuration"
|
||||
);
|
||||
let mut ordinary_marker = Vec::new();
|
||||
ordinary_archive
|
||||
.by_name(&format!("{UNREADABLE}/{EXPORT_UNREADABLE_MANIFEST}"))
|
||||
.expect("the ordinary export must name the configuration it could not read")
|
||||
.read_to_end(&mut ordinary_marker)
|
||||
.expect("read unreadable marker");
|
||||
assert!(
|
||||
!ordinary_marker
|
||||
.windows(SECRET.len())
|
||||
.any(|window| window == SECRET.as_bytes())
|
||||
);
|
||||
let ordinary_marker: serde_json::Value = serde_json::from_slice(&ordinary_marker).expect("the marker must be JSON");
|
||||
assert_eq!(ordinary_marker["bucket"], UNREADABLE);
|
||||
assert_eq!(
|
||||
ordinary_marker["unreadable"].as_array().map(Vec::len),
|
||||
Some(1),
|
||||
"only the configuration that could not be read may be marked: {ordinary_marker}"
|
||||
);
|
||||
assert_eq!(ordinary_marker["unreadable"][0]["config"], BUCKET_TARGETS_FILE);
|
||||
assert!(
|
||||
ordinary_marker["unreadable"][0]["error"].is_string(),
|
||||
"the marker must carry the reason an operator needs to repair the bucket"
|
||||
);
|
||||
|
||||
let response = ExportBucketMetadata {}
|
||||
.call(
|
||||
@@ -1676,6 +1841,12 @@ mod backup_zip_compatibility_tests {
|
||||
assert!(archive.by_name(DIAGNOSTIC_EXPORT_MANIFEST).is_err());
|
||||
assert!(archive.by_name(&format!("{HEALTHY}/{BUCKET_VERSIONING_CONFIG}")).is_ok());
|
||||
assert!(archive.by_name(&format!("{UNREADABLE}/{BUCKET_TARGETS_FILE}")).is_ok());
|
||||
assert!(
|
||||
archive
|
||||
.by_name(&format!("{UNREADABLE}/{EXPORT_UNREADABLE_MANIFEST}"))
|
||||
.is_err(),
|
||||
"the marker must disappear once the configuration reads again"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -1818,6 +1989,89 @@ mod backup_zip_compatibility_tests {
|
||||
}
|
||||
let _: ReplicationConfiguration =
|
||||
deserialize(&restored.replication_config_xml).expect("old parser must read the newly exported archive payload");
|
||||
|
||||
// rustfs/backlog#2309: a MinIO-origin `.metadata.bin` stores its
|
||||
// targets as a bare JSON array, which `BucketTargets` cannot decode.
|
||||
// Since rustfs/rustfs#7172 that reads as "stored but unreadable" —
|
||||
// which must mark one bucket's one configuration, not abort the
|
||||
// whole-cluster export an operator needs to migrate away.
|
||||
env.make_bucket(UNREADABLE_BUCKET, false).await;
|
||||
metadata_sys::update(UNREADABLE_BUCKET, BUCKET_TARGETS_FILE, MINIO_ARRAY_TARGETS.to_vec())
|
||||
.await
|
||||
.expect("persist the MinIO-shaped targets blob");
|
||||
metadata_sys::get_bucket_targets_config(UNREADABLE_BUCKET)
|
||||
.await
|
||||
.expect_err("a MinIO array-shaped targets blob must read as unreadable, not as an empty set");
|
||||
|
||||
let cluster_export = ExportBucketMetadata {}
|
||||
.call(
|
||||
admin_request(Method::GET, Uri::from_static("/rustfs/admin/v3/export-bucket-metadata"), Vec::new()),
|
||||
Params::new(),
|
||||
)
|
||||
.await
|
||||
.expect("one bucket's unreadable configuration must not abort the whole-cluster export");
|
||||
assert_eq!(cluster_export.output.0, StatusCode::OK);
|
||||
let cluster_archive = cluster_export
|
||||
.output
|
||||
.1
|
||||
.collect()
|
||||
.await
|
||||
.expect("read cluster archive body")
|
||||
.to_bytes()
|
||||
.to_vec();
|
||||
let mut archive = ZipArchive::new(Cursor::new(&cluster_archive)).expect("open cluster archive");
|
||||
|
||||
// Every readable configuration of every other bucket still exports.
|
||||
for (config_file, payload) in persisted_xml_fixtures() {
|
||||
let mut exported_payload = Vec::new();
|
||||
archive
|
||||
.by_name(&format!("{BUCKET}/{config_file}"))
|
||||
.unwrap_or_else(|_| panic!("cluster export must still contain {config_file}"))
|
||||
.read_to_end(&mut exported_payload)
|
||||
.unwrap_or_else(|_| panic!("read exported {config_file}"));
|
||||
assert_eq!(exported_payload, payload, "one bad bucket must not change another bucket's export");
|
||||
}
|
||||
assert!(
|
||||
archive.by_name(&format!("{BUCKET}/{EXPORT_UNREADABLE_MANIFEST}")).is_err(),
|
||||
"a bucket whose configurations all read must carry no unreadable marker"
|
||||
);
|
||||
|
||||
// The unreadable configuration is named rather than fabricated: no
|
||||
// targets entry is exported for it at all.
|
||||
assert!(
|
||||
archive
|
||||
.by_name(&format!("{UNREADABLE_BUCKET}/{BUCKET_TARGETS_FILE}"))
|
||||
.is_err(),
|
||||
"an unreadable targets blob must never be exported as a configuration"
|
||||
);
|
||||
let mut marker = Vec::new();
|
||||
archive
|
||||
.by_name(&format!("{UNREADABLE_BUCKET}/{EXPORT_UNREADABLE_MANIFEST}"))
|
||||
.expect("the export must name the configuration it could not read")
|
||||
.read_to_end(&mut marker)
|
||||
.expect("read unreadable marker");
|
||||
let marker: serde_json::Value = serde_json::from_slice(&marker).expect("the marker must be JSON");
|
||||
assert_eq!(marker["bucket"], UNREADABLE_BUCKET);
|
||||
assert_eq!(
|
||||
marker["unreadable"].as_array().map(Vec::len),
|
||||
Some(1),
|
||||
"only the configuration that could not be read may be marked: {marker}"
|
||||
);
|
||||
assert_eq!(marker["unreadable"][0]["config"], BUCKET_TARGETS_FILE);
|
||||
drop(archive);
|
||||
|
||||
// The marker cannot be misread as a configuration on the way back in:
|
||||
// the importer switches on configuration names and ignores everything
|
||||
// else, so the bucket's stored bytes come through untouched and the
|
||||
// operator still has to repair them explicitly.
|
||||
import_archive(cluster_archive).await;
|
||||
let after_round_trip = metadata_sys::get_config_from_disk(UNREADABLE_BUCKET)
|
||||
.await
|
||||
.expect("the marked bucket must still load after the round trip");
|
||||
assert_eq!(
|
||||
after_round_trip.bucket_targets_config_json, MINIO_ARRAY_TARGETS,
|
||||
"importing the marker must not overwrite or fabricate the bucket's targets configuration"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -18,20 +18,18 @@ use crate::admin::runtime_sources::object_store_from_extensions;
|
||||
use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket;
|
||||
use crate::admin::storage_api::error::StorageError;
|
||||
use crate::admin::storage_api::lifecycle::{
|
||||
IlmRecoveryClassification, IlmRecoveryProtocol, ManualTransitionCancelCheck, ManualTransitionJobRecord,
|
||||
ManualTransitionJobState, ManualTransitionProgressSink, ManualTransitionQueueSnapshot, ManualTransitionRunOptions,
|
||||
ManualTransitionRunReport, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim,
|
||||
TransitionOperatorDeleteResult, TransitionOperatorError, claim_manual_transition_scope_admission,
|
||||
delete_manual_transition_scope_admission_if_current, delete_transition_candidate_for_operator,
|
||||
enqueue_transition_for_existing_objects_scoped, finalize_missing_transition_transaction_for_operator,
|
||||
inspect_recovery_control, inspect_transition_transaction_for_operator, list_recovery_controls,
|
||||
ManualTransitionCancelCheck, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionProgressSink,
|
||||
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport, ManualTransitionScopeAdmission,
|
||||
ManualTransitionScopeAdmissionClaim, TransitionOperatorDeleteResult, TransitionOperatorError,
|
||||
claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
|
||||
delete_transition_candidate_for_operator, enqueue_transition_for_existing_objects_scoped,
|
||||
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
|
||||
load_manual_transition_job_record, load_manual_transition_scope_admission, manual_transition_job_lease_expired,
|
||||
manual_transition_queue_snapshot, manual_transition_scope_admission_lease_expired,
|
||||
persist_manual_transition_job_progress_if_owned, renew_manual_transition_job_lease_if_owned,
|
||||
request_manual_transition_job_cancel, save_manual_transition_job_record, update_manual_transition_job_record,
|
||||
};
|
||||
use crate::admin::storage_api::runtime::ECStore;
|
||||
use crate::admin::storage_api::s3::{S3ErrorCode as AdminS3ErrorCode, error as admin_s3_error};
|
||||
use crate::admin::utils::json_response;
|
||||
use crate::server::{ADMIN_PREFIX, RemoteAddr};
|
||||
use http::HeaderMap;
|
||||
@@ -232,48 +230,9 @@ pub fn register_ilm_transition_route(r: &mut S3Router<AdminOperation>) -> std::i
|
||||
format!("{ADMIN_PREFIX}/v3/ilm/transition/reconcile/{{transaction_id}}").as_str(),
|
||||
AdminOperation(&TransitionReconcileApplyHandler {}),
|
||||
)?;
|
||||
r.insert(
|
||||
Method::GET,
|
||||
format!("{ADMIN_PREFIX}/v3/ilm/recovery/records").as_str(),
|
||||
AdminOperation(&IlmRecoveryControlListHandler {}),
|
||||
)?;
|
||||
r.insert(
|
||||
Method::GET,
|
||||
format!("{ADMIN_PREFIX}/v3/ilm/recovery/records/{{control_id}}").as_str(),
|
||||
AdminOperation(&IlmRecoveryControlInspectHandler {}),
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct IlmRecoveryControlListQuery {
|
||||
protocol: IlmRecoveryProtocol,
|
||||
#[serde(default)]
|
||||
classification: Option<IlmRecoveryClassification>,
|
||||
#[serde(default = "default_recovery_control_list_limit")]
|
||||
limit: usize,
|
||||
#[serde(default)]
|
||||
marker: Option<String>,
|
||||
}
|
||||
|
||||
const fn default_recovery_control_list_limit() -> usize {
|
||||
100
|
||||
}
|
||||
|
||||
fn parse_recovery_control_list_query(query: Option<&str>) -> S3Result<IlmRecoveryControlListQuery> {
|
||||
let query = query.ok_or_else(|| admin_s3_error(AdminS3ErrorCode::InvalidRequest, "protocol is required"))?;
|
||||
let parsed: IlmRecoveryControlListQuery = serde_urlencoded::from_bytes(query.as_bytes())
|
||||
.map_err(|_| admin_s3_error(AdminS3ErrorCode::InvalidArgument, "invalid ILM recovery control query"))?;
|
||||
if !(1..=1_000).contains(&parsed.limit) {
|
||||
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "limit must be between 1 and 1000"));
|
||||
}
|
||||
if parsed.marker.as_ref().is_some_and(|marker| marker.is_empty()) {
|
||||
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "marker must not be empty"));
|
||||
}
|
||||
Ok(parsed)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum ManualTransitionRunMode {
|
||||
EnqueueOnly,
|
||||
@@ -464,26 +423,6 @@ fn transition_transaction_id_from_params(params: &Params<'_, '_>) -> S3Result<Uu
|
||||
.map_err(|_| s3_error!(InvalidArgument, "invalid transition transaction id"))
|
||||
}
|
||||
|
||||
fn recovery_control_id_from_params(params: &Params<'_, '_>) -> S3Result<String> {
|
||||
let control_id = params.get("control_id").unwrap_or("");
|
||||
if control_id.len() != 64
|
||||
|| !control_id
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
|
||||
{
|
||||
return Err(admin_s3_error(AdminS3ErrorCode::InvalidArgument, "invalid ILM recovery control id"));
|
||||
}
|
||||
Ok(control_id.to_string())
|
||||
}
|
||||
|
||||
fn map_recovery_control_error(err: StorageError) -> S3Error {
|
||||
if err == StorageError::ConfigNotFound {
|
||||
admin_s3_error(AdminS3ErrorCode::NoSuchKey, "ILM recovery control not found")
|
||||
} else {
|
||||
admin_s3_error(AdminS3ErrorCode::InternalError, "ILM recovery control request failed")
|
||||
}
|
||||
}
|
||||
|
||||
fn map_transition_operator_error(err: TransitionOperatorError) -> S3Error {
|
||||
match err {
|
||||
TransitionOperatorError::NotFound => s3_error!(NoSuchKey, "transition transaction not found"),
|
||||
@@ -1092,40 +1031,6 @@ impl Operation for TransitionReconcileInspectHandler {
|
||||
}
|
||||
}
|
||||
|
||||
pub struct IlmRecoveryControlListHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for IlmRecoveryControlListHandler {
|
||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
authorize_transition_admin_request(&req, AdminAction::ListTierAction).await?;
|
||||
let query = parse_recovery_control_list_query(req.uri.query())?;
|
||||
let Some(store) = object_store_from_extensions(&req.extensions) else {
|
||||
return Err(admin_s3_error(AdminS3ErrorCode::InternalError, "object store is not initialized"));
|
||||
};
|
||||
let page = list_recovery_controls(store, query.protocol, query.classification, query.limit, query.marker)
|
||||
.await
|
||||
.map_err(map_recovery_control_error)?;
|
||||
json_response(StatusCode::OK, &page)
|
||||
}
|
||||
}
|
||||
|
||||
pub struct IlmRecoveryControlInspectHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for IlmRecoveryControlInspectHandler {
|
||||
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
authorize_transition_admin_request(&req, AdminAction::ListTierAction).await?;
|
||||
let control_id = recovery_control_id_from_params(¶ms)?;
|
||||
let Some(store) = object_store_from_extensions(&req.extensions) else {
|
||||
return Err(admin_s3_error(AdminS3ErrorCode::InternalError, "object store is not initialized"));
|
||||
};
|
||||
let control = inspect_recovery_control(store, &control_id)
|
||||
.await
|
||||
.map_err(map_recovery_control_error)?;
|
||||
json_response(StatusCode::OK, &control)
|
||||
}
|
||||
}
|
||||
|
||||
pub struct TransitionReconcileApplyHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
@@ -1199,49 +1104,6 @@ mod tests {
|
||||
f(&matched.params)
|
||||
}
|
||||
|
||||
fn with_recovery_control_params<T>(path: &str, f: impl FnOnce(&Params<'_, '_>) -> T) -> T {
|
||||
let mut router = Router::new();
|
||||
router
|
||||
.insert("/rustfs/admin/v3/ilm/recovery/records/{control_id}", ())
|
||||
.expect("route should insert");
|
||||
let matched = router.at(path).expect("route should match");
|
||||
f(&matched.params)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recovery_control_query_is_bounded_and_strict() {
|
||||
let query = parse_recovery_control_list_query(Some("protocol=transition_transaction"))
|
||||
.expect("minimal recovery query should parse");
|
||||
assert_eq!(query.protocol, IlmRecoveryProtocol::TransitionTransaction);
|
||||
assert_eq!(query.classification, None);
|
||||
assert_eq!(query.limit, 100);
|
||||
|
||||
let filtered = parse_recovery_control_list_query(Some(
|
||||
"protocol=tier_delete_journal&classification=retained_ambiguous&limit=1000&marker=opaque",
|
||||
))
|
||||
.expect("bounded filtered query should parse");
|
||||
assert_eq!(filtered.protocol, IlmRecoveryProtocol::TierDeleteJournal);
|
||||
assert_eq!(filtered.classification, Some(IlmRecoveryClassification::RetainedAmbiguous));
|
||||
assert_eq!(filtered.limit, 1000);
|
||||
assert!(parse_recovery_control_list_query(None).is_err());
|
||||
assert!(parse_recovery_control_list_query(Some("protocol=transition_transaction&limit=0")).is_err());
|
||||
assert!(parse_recovery_control_list_query(Some("protocol=transition_transaction&limit=1001")).is_err());
|
||||
assert!(parse_recovery_control_list_query(Some("protocol=unknown")).is_err());
|
||||
assert!(parse_recovery_control_list_query(Some("protocol=transition_transaction&extra=true")).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recovery_control_id_is_canonical_lowercase_sha256() {
|
||||
let id = "ab".repeat(32);
|
||||
with_recovery_control_params(&format!("/rustfs/admin/v3/ilm/recovery/records/{id}"), |params| {
|
||||
assert_eq!(recovery_control_id_from_params(params).expect("control id should parse"), id);
|
||||
});
|
||||
let uppercase = "AB".repeat(32);
|
||||
with_recovery_control_params(&format!("/rustfs/admin/v3/ilm/recovery/records/{uppercase}"), |params| {
|
||||
assert!(recovery_control_id_from_params(params).is_err())
|
||||
});
|
||||
}
|
||||
|
||||
fn manual_transition_job_request(method: Method, path: &'static str) -> S3Request<Body> {
|
||||
S3Request {
|
||||
input: Body::empty(),
|
||||
|
||||
@@ -29,7 +29,7 @@ use crate::admin::storage_api::bucket::replication::{REMOTE_TARGET_READ_ONLY_HIS
|
||||
use crate::admin::storage_api::bucket::target::{
|
||||
BucketTarget, BucketTargetType, Credentials as TargetCredentials, LatencyStat, duration_from_secs_or_nanos,
|
||||
};
|
||||
use crate::admin::storage_api::bucket::target_sys::{BucketTargetError, BucketTargetSys};
|
||||
use crate::admin::storage_api::bucket::target_sys::{BucketTargetError, BucketTargetSys, UnreadableTargetsPolicy};
|
||||
use crate::admin::storage_api::contract::bucket::{BucketOperations, BucketOptions};
|
||||
use crate::admin::storage_api::contract::list::ListOperations as _;
|
||||
use crate::admin::storage_api::error::StorageError;
|
||||
@@ -57,6 +57,20 @@ use url::Host;
|
||||
|
||||
const SUPPORTED_REMOTE_TARGET_API: &str = "s3v4";
|
||||
|
||||
const LOG_COMPONENT_ADMIN: &str = "admin";
|
||||
const LOG_SUBSYSTEM_REPLICATION: &str = "replication";
|
||||
const EVENT_ADMIN_REMOTE_TARGET_STATE: &str = "admin_remote_target_state";
|
||||
|
||||
/// `set-remote-target?replace-unreadable=true`: the operator's explicit
|
||||
/// acknowledgement that this bucket's persisted target set cannot be decoded
|
||||
/// and is to be discarded (rustfs/backlog#2309).
|
||||
///
|
||||
/// Without it the refusal from rustfs/rustfs#7172 stands, which is what keeps
|
||||
/// an unreadable set from being silently rewritten from a partial view. The
|
||||
/// flag is deliberately absent from `PutBucketReplication`: targets are
|
||||
/// repaired first, then the rule is set.
|
||||
const REPLACE_UNREADABLE_TARGETS_PARAM: &str = "replace-unreadable";
|
||||
|
||||
/// Field groups a `set-remote-target?update=true` request may modify, mirroring
|
||||
/// MinIO's `TargetUpdateType` / `GetTargetUpdateOps` query contract: the update
|
||||
/// overlays only the requested groups onto the stored target, so a client can
|
||||
@@ -541,6 +555,7 @@ impl Operation for SetRemoteTargetHandler {
|
||||
};
|
||||
|
||||
let update = queries.get("update").is_some_and(|v| v == "true");
|
||||
let replace_unreadable = queries.get(REPLACE_UNREADABLE_TARGETS_PARAM).is_some_and(|v| v == "true");
|
||||
|
||||
warn!("set remote target, bucket: {}, update: {}", bucket, update);
|
||||
|
||||
@@ -708,10 +723,37 @@ impl Operation for SetRemoteTargetHandler {
|
||||
|
||||
let arn = remote_target.arn.clone();
|
||||
|
||||
let unreadable_policy = if replace_unreadable {
|
||||
UnreadableTargetsPolicy::Replace
|
||||
} else {
|
||||
UnreadableTargetsPolicy::FailClosed
|
||||
};
|
||||
let discarding_unreadable_targets = replace_unreadable
|
||||
&& matches!(
|
||||
bucket_target_sys.list_bucket_targets(bucket).await,
|
||||
Err(BucketTargetError::BucketRemoteTargetsUnreadable { .. })
|
||||
);
|
||||
|
||||
let targets = bucket_target_sys
|
||||
.set_target(bucket, &remote_target, update)
|
||||
.set_target(bucket, &remote_target, update, unreadable_policy)
|
||||
.await
|
||||
.map_err(map_bucket_target_error)?;
|
||||
|
||||
// Audited only where the discard actually happened: the flag alone is
|
||||
// not an event, on a readable set it changes nothing, and a refused
|
||||
// write must not leave a record claiming the set was replaced.
|
||||
if discarding_unreadable_targets {
|
||||
warn!(
|
||||
event = EVENT_ADMIN_REMOTE_TARGET_STATE,
|
||||
component = LOG_COMPONENT_ADMIN,
|
||||
subsystem = LOG_SUBSYSTEM_REPLICATION,
|
||||
action = "set_remote_target",
|
||||
result = "unreadable_targets_replaced",
|
||||
bucket = %bucket,
|
||||
arn = %remote_target.arn,
|
||||
"admin remote target state"
|
||||
);
|
||||
}
|
||||
let json_targets = serde_json::to_vec(&targets).map_err(|e| {
|
||||
error!("Serialization error: {}", e);
|
||||
S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize targets".to_string())
|
||||
|
||||
@@ -212,6 +212,7 @@ pub(crate) mod bucket_target_sys {
|
||||
pub(crate) type S3ClientError = super::ecstore_bucket::bucket_target_sys::S3ClientError;
|
||||
pub(crate) type SsecPassthroughCapability = super::ecstore_bucket::bucket_target_sys::SsecPassthroughCapability;
|
||||
pub(crate) type TargetClient = super::ecstore_bucket::bucket_target_sys::TargetClient;
|
||||
pub(crate) type UnreadableTargetsPolicy = super::ecstore_bucket::bucket_target_sys::UnreadableTargetsPolicy;
|
||||
}
|
||||
|
||||
pub(crate) mod lifecycle {
|
||||
@@ -232,9 +233,6 @@ pub(crate) mod lifecycle {
|
||||
pub(crate) type ManualTransitionRunOptions =
|
||||
super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunOptions;
|
||||
pub(crate) type ManualTransitionRunReport = super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunReport;
|
||||
pub(crate) use super::ecstore_bucket::lifecycle::recovery_control::{
|
||||
IlmRecoveryClassification, IlmRecoveryProtocol, inspect_recovery_control, list_recovery_controls,
|
||||
};
|
||||
pub(crate) use super::ecstore_bucket::lifecycle::transition_transaction::{
|
||||
TransitionOperatorDeleteResult, TransitionOperatorError, delete_transition_candidate_for_operator,
|
||||
finalize_missing_transition_transaction_for_operator, inspect_transition_transaction_for_operator,
|
||||
|
||||
@@ -133,7 +133,6 @@ impl NativeHttp {
|
||||
self.send_classified(request, error_code_header, false).await
|
||||
}
|
||||
|
||||
#[cfg(feature = "gcs")]
|
||||
pub(super) async fn send_object(
|
||||
&self,
|
||||
request: reqwest::Request,
|
||||
@@ -211,7 +210,6 @@ pub(super) async fn read_text(response: reqwest::Response, max_bytes: usize) ->
|
||||
/// Base64 digest (`Content-MD5`, `md5Hash`, `x-goog-hash`) as lowercase hex.
|
||||
/// `None` when the value is not a 16-byte digest, so a CRC32C never passes as
|
||||
/// an MD5.
|
||||
#[cfg(any(test, feature = "gcs"))]
|
||||
pub(super) fn base64_md5_to_hex(value: &str) -> Option<String> {
|
||||
let raw = base64_simd::STANDARD.decode_to_vec(value.trim().as_bytes()).ok()?;
|
||||
(raw.len() == 16).then(|| faster_hex::hex_string(&raw))
|
||||
|
||||
@@ -4034,7 +4034,7 @@ mod tests {
|
||||
let enabled_before = sys.is_module_enabled();
|
||||
sys.set_module_enabled(true);
|
||||
let bucket = format!("odm-capture-failure-{}", uuid::Uuid::new_v4());
|
||||
let mut config: crate::on_demand_migration::OnDemandMigrationConfig = serde_json::from_str(r#"{"source":{"provider":"minio","endpoint":"https://source.example.com","region":"us-east-1","bucket":"source","credentials":{"access_key":"test","secret_key":"test"}}}"#).expect("source config");
|
||||
let mut config: crate::on_demand_migration::OnDemandMigrationConfig = serde_json::from_str(r#"{"source":{"provider":"minio","endpoint":"https://source.example.com","bucket":"source","credentials":{"access_key":"test","secret_key":"test"}}}"#).expect("source config");
|
||||
config.policy.list_through = true;
|
||||
sys.apply_for_incarnation(&bucket, uuid::Uuid::new_v4(), Some(&config)).await;
|
||||
let mut get = build_request(
|
||||
|
||||
@@ -50,8 +50,13 @@ cd "$(dirname "$0")/.."
|
||||
# s3_error! stays flat at 1616.
|
||||
# 1616 → 1613 on 2026-09-02: dependency refresh verified the current tree has
|
||||
# already shed three s3_error! invocation lines; retighten the line counter.
|
||||
# 1613 → 1589 on 2026-09-06: rustfs/backlog#2309 and rustfs/rustfs#7225 both
|
||||
# folded the ten per-config arms of ExportBucketMetadata into one helper, which
|
||||
# now reports an unreadable configuration as a plain string instead of raising
|
||||
# an S3 error per arm (24 invocation lines removed from
|
||||
# rustfs/src/admin/handlers/bucket_meta.rs; measured after merging the two).
|
||||
S3S_IMPORT_FILES_BASELINE=213
|
||||
S3_ERROR_LINES_BASELINE=1613
|
||||
S3_ERROR_LINES_BASELINE=1589
|
||||
# ecstore-scoped ratchet (rustfs/backlog#1842): the storage engine must not
|
||||
# know S3 wire/DTO types (ARCHITECTURE.md invariant 4). The S3-*consuming*
|
||||
# client was extracted to crates/s3-client, where s3s usage is legitimate;
|
||||
|
||||
Reference in New Issue
Block a user