fix(replication): resolve drifted replicas via a target version ledger

A replication target that mints its own version ids (Wasabi, AWS S3)
never answers to the source uuid, so every version-addressed mutation
after the initial PUT failed forever: permanent version deletes answered
NoSuchVersion every heal cycle, and tag / retention / legal-hold updates
re-PUT the object, minting one more target version per update
(rustfs/backlog#2340).

Record the id the target assigned as a per-target ledger on the source
version (replication-target-version-<arn>, written through the existing
status writeback) and resolve every later mutation through it: version
deletes DELETE the ledger id, metadata updates go through the
metadata-only Object Lock and tagging APIs. Replicas written before the
ledger existed are located by exact key and ETag, minus the candidates
other generations of the key already claim through their own ledgers; an
ambiguous remainder is refused with a backoff instead of guessed, since a
wrong pick would destroy a live generation. A fresh write never consults
content identity. NoSuchVersion on a version-addressed DELETE counts as
purged.

The fake target gains the Wasabi shape (404 NoSuchVersion on an unknown
id, per-version Object Lock APIs) and the matrix covers the three
mutation classes plus the same-bytes generation case.
This commit is contained in:
唐小鸭
2026-09-07 15:39:44 +08:00
parent 481c1b7939
commit 46a387dffe
11 changed files with 2117 additions and 224 deletions
+207 -5
View File
@@ -34,12 +34,14 @@ use s3s::dto::{
AbortMultipartUploadInput, AbortMultipartUploadOutput, CommonPrefix, CompleteMultipartUploadInput,
CompleteMultipartUploadOutput, CreateMultipartUploadInput, CreateMultipartUploadOutput, DeleteMarkerEntry, DeleteObjectInput,
DeleteObjectOutput, DeleteObjectTaggingInput, DeleteObjectTaggingOutput, ETag, GetBucketVersioningInput,
GetBucketVersioningOutput, GetObjectInput, GetObjectLockConfigurationInput, GetObjectLockConfigurationOutput,
GetObjectOutput, GetObjectTaggingInput, GetObjectTaggingOutput, HeadBucketInput, HeadBucketOutput, HeadObjectInput,
GetBucketVersioningOutput, GetObjectInput, GetObjectLegalHoldInput, GetObjectLegalHoldOutput,
GetObjectLockConfigurationInput, GetObjectLockConfigurationOutput, GetObjectOutput, GetObjectRetentionInput,
GetObjectRetentionOutput, GetObjectTaggingInput, GetObjectTaggingOutput, HeadBucketInput, HeadBucketOutput, HeadObjectInput,
HeadObjectOutput, ListObjectVersionsInput, ListObjectVersionsOutput, ListObjectsV2Input, ListObjectsV2Output, Object,
ObjectLockConfiguration, ObjectLockEnabled, ObjectStorageClass, ObjectVersionId, PutObjectInput, PutObjectOutput,
PutObjectTaggingInput, PutObjectTaggingOutput, Range, StreamingBlob, Tag, TagSet, Timestamp, TimestampFormat,
UploadPartInput, UploadPartOutput,
ObjectLockConfiguration, ObjectLockEnabled, ObjectLockLegalHold, ObjectLockLegalHoldStatus, ObjectLockMode,
ObjectLockRetention, ObjectLockRetentionMode, ObjectStorageClass, ObjectVersionId, PutObjectInput, PutObjectLegalHoldInput,
PutObjectLegalHoldOutput, PutObjectOutput, PutObjectRetentionInput, PutObjectRetentionOutput, PutObjectTaggingInput,
PutObjectTaggingOutput, Range, StreamingBlob, Tag, TagSet, Timestamp, TimestampFormat, UploadPartInput, UploadPartOutput,
};
use s3s::service::{S3Service, S3ServiceBuilder};
use s3s::validation::{AwsNameValidation, NameValidation};
@@ -127,6 +129,10 @@ pub enum Operation {
GetObjectTagging,
PutObjectTagging,
DeleteObjectTagging,
GetObjectRetention,
PutObjectRetention,
GetObjectLegalHold,
PutObjectLegalHold,
ListObjectVersions,
ListObjectsV2,
CreateMultipartUpload,
@@ -501,6 +507,10 @@ struct StoreState {
/// PutObject carrying any `x-amz-object-lock-*` header must also carry
/// `Content-MD5` or an `x-amz-checksum-*` header.
require_checksum_for_object_lock: bool,
/// Models Wasabi (rustfs/backlog#2340): a version-addressed DELETE of a
/// version id the target never had answers 404 `NoSuchVersion` instead of
/// the idempotent 204 RustFS/MinIO give.
reject_unknown_version_deletes: bool,
limits: StoreLimits,
buckets: HashMap<String, BucketState>,
uploads: HashMap<String, MultipartState>,
@@ -565,6 +575,41 @@ struct ObjectVersion {
/// SSE-C passthrough transport headers stored with the version (RustFS
/// target behavior); empty when the drop mode discarded them.
replication_sse_headers: Vec<(String, String)>,
/// Object Lock state of the version: retention (mode, retain-until) from
/// the PUT / CreateMultipartUpload headers or PutObjectRetention, and the
/// legal hold flag; replayed on HEAD.
lock: VersionLock,
}
#[derive(Clone, Default)]
struct VersionLock {
retention: Option<(String, Timestamp)>,
/// `None` until a legal hold status was ever set; like S3, HEAD then
/// reports nothing, while an explicit OFF is reported as `OFF`.
legal_hold: Option<bool>,
}
impl VersionLock {
fn from_headers(
mode: Option<ObjectLockMode>,
retain_until: Option<Timestamp>,
legal_hold: Option<ObjectLockLegalHoldStatus>,
) -> Self {
Self {
retention: mode.zip(retain_until).map(|(mode, until)| (mode.as_str().to_string(), until)),
legal_hold: legal_hold.map(|status| status.as_str().eq_ignore_ascii_case("ON")),
}
}
fn legal_hold_status(&self) -> Option<ObjectLockLegalHoldStatus> {
self.legal_hold.map(|on| {
ObjectLockLegalHoldStatus::from_static(if on {
ObjectLockLegalHoldStatus::ON
} else {
ObjectLockLegalHoldStatus::OFF
})
})
}
}
#[derive(Clone)]
@@ -576,6 +621,7 @@ struct MultipartState {
metadata: Option<HashMap<String, String>>,
standard_headers: StandardHeaders,
replication_sse_headers: Vec<(String, String)>,
lock: VersionLock,
parts: BTreeMap<i32, MultipartPart>,
}
@@ -845,6 +891,7 @@ impl FakeS3Target {
standard_headers: seed.standard_headers.clone(),
tags: Vec::new(),
replication_sse_headers: Vec::new(),
lock: VersionLock::default(),
};
upsert_version(&mut state, bucket, key.into(), version).expect("seed object must fit the storage budget");
e_tag
@@ -918,6 +965,12 @@ impl FakeS3Target {
/// PutObject that carries Object Lock parameters (AWS S3 / MinIO rule,
/// rustfs#7082). `Content-MD5`, when present, is always verified against
/// the body regardless of this mode.
/// Wasabi-like mode: DELETE of an unknown version id answers 404
/// `NoSuchVersion` (the default 204 models RustFS/MinIO).
pub fn reject_unknown_version_deletes(&self, enabled: bool) {
lock(&self.backend.store).reject_unknown_version_deletes = enabled;
}
pub fn require_checksum_for_object_lock(&self, enabled: bool) {
lock(&self.backend.store).require_checksum_for_object_lock = enabled;
}
@@ -1152,6 +1205,10 @@ fn operation_from_s3_name(name: &str) -> Operation {
"GetObjectTagging" => Operation::GetObjectTagging,
"PutObjectTagging" => Operation::PutObjectTagging,
"DeleteObjectTagging" => Operation::DeleteObjectTagging,
"GetObjectRetention" => Operation::GetObjectRetention,
"PutObjectRetention" => Operation::PutObjectRetention,
"GetObjectLegalHold" => Operation::GetObjectLegalHold,
"PutObjectLegalHold" => Operation::PutObjectLegalHold,
"ListObjectsV2" => Operation::ListObjectsV2,
"CreateMultipartUpload" => Operation::CreateMultipartUpload,
"UploadPart" => Operation::UploadPart,
@@ -1289,6 +1346,18 @@ fn parse_request(method: &Method, uri: &Uri) -> ParsedRequest {
(&Method::DELETE, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::DeleteObjectTagging
}
(&Method::GET, true) if query.contains_key("retention") && only_query_keys(&["retention", "versionId"]) => {
Operation::GetObjectRetention
}
(&Method::PUT, true) if query.contains_key("retention") && only_query_keys(&["retention", "versionId"]) => {
Operation::PutObjectRetention
}
(&Method::GET, true) if query.contains_key("legal-hold") && only_query_keys(&["legal-hold", "versionId"]) => {
Operation::GetObjectLegalHold
}
(&Method::PUT, true) if query.contains_key("legal-hold") && only_query_keys(&["legal-hold", "versionId"]) => {
Operation::PutObjectLegalHold
}
// A replication PUT addresses the source version via `?versionId=`.
(&Method::PUT, true) if only_query_keys(&["versionId"]) => Operation::PutObject,
(&Method::GET, true) if only_query_keys(&["versionId"]) => Operation::GetObject,
@@ -1842,6 +1911,28 @@ fn set_version_tags(
Ok(resolved)
}
fn update_version_lock(
state: &mut StoreState,
bucket: &str,
key: &str,
version_id: Option<&str>,
update: impl FnOnce(&mut VersionLock),
) -> S3Result<String> {
let resolved = find_version(state, bucket, key, version_id)?.version_id;
let version = state
.buckets
.get_mut(bucket)
.expect("bucket existence checked by find_version")
.objects
.get_mut(key)
.expect("key existence checked by find_version")
.iter_mut()
.find(|version| version.version_id == resolved)
.expect("version existence checked by find_version");
update(&mut version.lock);
Ok(resolved)
}
/// Whether version ids are surfaced for this bucket. Unknown buckets report
/// `true`; the caller's lookup raises `NoSuchBucket` first.
fn bucket_versioned(state: &StoreState, bucket: &str) -> bool {
@@ -2281,6 +2372,11 @@ impl S3 for FakeBackend {
standard_headers,
tags: Vec::new(),
replication_sse_headers: captured_replication_sse_headers(&headers, drop_unlisted),
lock: VersionLock::from_headers(
input.object_lock_mode,
input.object_lock_retain_until_date,
input.object_lock_legal_hold_status,
),
};
upsert_version(&mut lock(&self.store), &input.bucket, input.key, version)?;
Ok(apply_response_fault(
@@ -2339,6 +2435,13 @@ impl S3 for FakeBackend {
last_modified: Some(version.last_modified.clone()),
version_id: versioned.then_some(version.version_id),
sse_customer_algorithm,
object_lock_mode: version
.lock
.retention
.as_ref()
.map(|(mode, _)| ObjectLockMode::from(mode.clone())),
object_lock_retain_until_date: version.lock.retention.as_ref().map(|(_, until)| until.clone()),
object_lock_legal_hold_status: version.lock.legal_hold_status(),
..Default::default()
});
response.status = served.status;
@@ -2373,6 +2476,13 @@ impl S3 for FakeBackend {
last_modified: Some(version.last_modified.clone()),
version_id: versioned.then_some(version.version_id),
sse_customer_algorithm,
object_lock_mode: version
.lock
.retention
.as_ref()
.map(|(mode, _)| ObjectLockMode::from(mode.clone())),
object_lock_retain_until_date: version.lock.retention.as_ref().map(|(_, until)| until.clone()),
object_lock_legal_hold_status: version.lock.legal_hold_status(),
..Default::default()
});
response.status = served.status;
@@ -2432,6 +2542,82 @@ impl S3 for FakeBackend {
))
}
async fn get_object_retention(
&self,
req: S3Request<GetObjectRetentionInput>,
) -> S3Result<S3Response<GetObjectRetentionOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let version = find_version(&lock(&self.store), &input.bucket, &input.key, input.version_id.as_deref())?;
Ok(apply_response_fault(
S3Response::new(GetObjectRetentionOutput {
retention: version.lock.retention.map(|(mode, until)| ObjectLockRetention {
mode: Some(ObjectLockRetentionMode::from(mode)),
retain_until_date: Some(until),
}),
}),
fault.as_ref(),
))
}
async fn put_object_retention(
&self,
req: S3Request<PutObjectRetentionInput>,
) -> S3Result<S3Response<PutObjectRetentionOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let retention = input
.retention
.and_then(|retention| retention.mode.zip(retention.retain_until_date))
.map(|(mode, until)| (mode.as_str().to_string(), until));
update_version_lock(&mut lock(&self.store), &input.bucket, &input.key, input.version_id.as_deref(), |lock| {
lock.retention = retention;
})?;
Ok(apply_response_fault(S3Response::new(PutObjectRetentionOutput::default()), fault.as_ref()))
}
async fn get_object_legal_hold(
&self,
req: S3Request<GetObjectLegalHoldInput>,
) -> S3Result<S3Response<GetObjectLegalHoldOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let version = find_version(&lock(&self.store), &input.bucket, &input.key, input.version_id.as_deref())?;
Ok(apply_response_fault(
S3Response::new(GetObjectLegalHoldOutput {
legal_hold: Some(ObjectLockLegalHold {
status: Some(
version
.lock
.legal_hold_status()
.unwrap_or_else(|| ObjectLockLegalHoldStatus::from_static(ObjectLockLegalHoldStatus::OFF)),
),
}),
}),
fault.as_ref(),
))
}
async fn put_object_legal_hold(
&self,
req: S3Request<PutObjectLegalHoldInput>,
) -> S3Result<S3Response<PutObjectLegalHoldOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let legal_hold_on = input
.legal_hold
.and_then(|hold| hold.status)
.is_some_and(|status| status.as_str().eq_ignore_ascii_case("ON"));
update_version_lock(&mut lock(&self.store), &input.bucket, &input.key, input.version_id.as_deref(), |lock| {
lock.legal_hold = Some(legal_hold_on);
})?;
Ok(apply_response_fault(S3Response::new(PutObjectLegalHoldOutput::default()), fault.as_ref()))
}
async fn delete_object_tagging(
&self,
req: S3Request<DeleteObjectTaggingInput>,
@@ -2485,6 +2671,7 @@ impl S3 for FakeBackend {
return Ok(apply_response_fault(S3Response::new(DeleteObjectOutput::default()), fault.as_ref()));
}
if let Some(version_id) = input.version_id {
let reject_unknown = state.reject_unknown_version_deletes;
let (removed_bytes, removed_versions, delete_marker, remove_key) = {
let Some(versions) = state
.buckets
@@ -2493,6 +2680,9 @@ impl S3 for FakeBackend {
.objects
.get_mut(&input.key)
else {
if reject_unknown {
return Err(s3s::s3_error!(NoSuchVersion, "The specified version does not exist."));
}
return Ok(apply_response_fault(
S3Response::new(DeleteObjectOutput {
version_id: Some(version_id),
@@ -2501,6 +2691,9 @@ impl S3 for FakeBackend {
fault.as_ref(),
));
};
if reject_unknown && !versions.iter().any(|version| version.version_id == version_id) {
return Err(s3s::s3_error!(NoSuchVersion, "The specified version does not exist."));
}
let mut removed_bytes = 0usize;
let mut removed_versions = 0usize;
let mut delete_marker = None;
@@ -2554,6 +2747,7 @@ impl S3 for FakeBackend {
standard_headers: StandardHeaders::default(),
tags: Vec::new(),
replication_sse_headers: Vec::new(),
lock: VersionLock::default(),
},
)?;
Ok(apply_response_fault(
@@ -2608,6 +2802,11 @@ impl S3 for FakeBackend {
metadata: input.metadata,
standard_headers,
replication_sse_headers: captured_replication_sse_headers(&headers, drop_unlisted),
lock: VersionLock::from_headers(
input.object_lock_mode,
input.object_lock_retain_until_date,
input.object_lock_legal_hold_status,
),
parts: BTreeMap::new(),
},
);
@@ -2748,6 +2947,7 @@ impl S3 for FakeBackend {
metadata: upload.metadata.clone(),
standard_headers: upload.standard_headers.clone(),
replication_sse_headers: upload.replication_sse_headers.clone(),
lock: upload.lock.clone(),
parts: BTreeMap::new(),
},
selected,
@@ -2778,6 +2978,7 @@ impl S3 for FakeBackend {
standard_headers: upload.standard_headers,
tags: Vec::new(),
replication_sse_headers: upload.replication_sse_headers,
lock: upload.lock,
};
let mut state = lock(&self.store);
let versioned = bucket_versioned(&state, &input.bucket);
@@ -4609,6 +4810,7 @@ mod tests {
metadata: None,
standard_headers: StandardHeaders::default(),
replication_sse_headers: Vec::new(),
lock: VersionLock::default(),
parts: BTreeMap::new(),
},
);
@@ -512,7 +512,7 @@ pub(crate) async fn put_bucket_replication(
put_bucket_replication_with_delete_statuses(env, bucket, target_arn, "Enabled", None).await
}
async fn put_bucket_replication_with_delete_statuses(
pub(crate) async fn put_bucket_replication_with_delete_statuses(
env: &RustFSTestEnvironment,
bucket: &str,
target_arn: &str,
@@ -37,13 +37,14 @@ use crate::fake_s3_target::{FakeS3Target, FaultAction as FakeTargetFault, Operat
use crate::on_demand_migration::common::{OdmEnvOptions, OdmTestEnv, fake_source_client};
use crate::replication_extension_test::{
LOOPBACK_REPLICATION_TARGET_ENV, ReplicationTargetOptions, enable_bucket_versioning, get_replication_reset_status,
put_bucket_replication, set_replication_target_with_options, start_bucket_replication_reset,
put_bucket_replication, put_bucket_replication_with_delete_statuses, set_replication_target_with_options,
start_bucket_replication_reset,
};
use aws_sdk_s3::Client;
use aws_sdk_s3::primitives::{ByteStream, DateTime};
use aws_sdk_s3::types::{
Checksum, ChecksumAlgorithm, CompletedMultipartUpload, CompletedPart, ObjectAttributes, ObjectLockLegalHoldStatus,
ObjectLockMode,
Checksum, ChecksumAlgorithm, CompletedMultipartUpload, CompletedPart, ObjectAttributes, ObjectLockLegalHold,
ObjectLockLegalHoldStatus, ObjectLockMode, ObjectLockRetention, ObjectLockRetentionMode, Tag, Tagging,
};
use bytes::Bytes;
use std::error::Error;
@@ -66,7 +67,10 @@ enum TargetMode {
/// Object Lock parameters must carry `Content-MD5` or `x-amz-checksum-*`.
RequireChecksumWithObjectLock,
/// AWS S3 / Wasabi / Impossible Cloud: mints its own version ids
/// (rustfs/backlog#2085). Data must still land.
/// (rustfs/backlog#2085) and, like Wasabi, answers NoSuchVersion to a
/// DELETE of an id it never had (rustfs/backlog#2340). Data must still
/// land, and every version-addressed mutation must resolve the replica
/// through the target-version ledger.
MintOwnVersionIds,
}
@@ -83,7 +87,10 @@ impl TargetMode {
TargetMode::Baseline => {}
TargetMode::RejectAwsChunked => target.reject_aws_chunked_uploads(true),
TargetMode::RequireChecksumWithObjectLock => target.require_checksum_for_object_lock(true),
TargetMode::MintOwnVersionIds => target.assign_own_version_ids(true),
TargetMode::MintOwnVersionIds => {
target.assign_own_version_ids(true);
target.reject_unknown_version_deletes(true);
}
}
}
@@ -384,6 +391,296 @@ async fn matrix_mint_own_version_ids_redrive_does_not_duplicate() -> TestResult
Ok(())
}
/// rustfs/backlog#2340 (target-version ledger): on a target that mints its own
/// version ids and answers NoSuchVersion to an unknown id (the Wasabi shape),
/// every version-addressed mutation must land on the version the target
/// assigned, which the replication PUT recorded on the source:
/// - a tag update changes the existing target version, no new version;
/// - a retention extension and legal hold ON/OFF change that version too;
/// - a permanent delete of the older of two same-content generations removes
/// exactly that replica and keeps the live one (content identity alone
/// could not tell them apart).
#[tokio::test]
async fn matrix_mint_own_version_ids_addresses_mutations_through_the_ledger() -> TestResult {
init_logging();
let target = FakeS3Target::start().await?;
let target_bucket = "matrix-mint-own-ledger-dst".to_string();
target.create_bucket_with_object_lock(target_bucket.clone());
TargetMode::MintOwnVersionIds.apply(&target);
let mut env_vars = replication_fast_env();
env_vars.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
env_vars.extend_from_slice(&[
("NO_PROXY", "127.0.0.1,localhost"),
("HTTP_PROXY", ""),
("HTTPS_PROXY", ""),
// The scanner heal pass retries a purge the first attempt lost.
("RUSTFS_SCANNER_CYCLE", "1"),
("RUSTFS_SCANNER_START_DELAY_SECS", "1"),
]);
let env = OdmTestEnv::start_with(OdmEnvOptions {
env: env_vars,
..OdmEnvOptions::default()
})
.await?;
let source_env = &env.rustfs;
let source_bucket = "matrix-mint-own-ledger-src";
let source_client = source_env.create_s3_client();
source_client
.create_bucket()
.bucket(source_bucket)
.object_lock_enabled_for_bucket(true)
.send()
.await?;
enable_bucket_versioning(source_env, source_bucket).await?;
let target_arn = set_replication_target_with_options(
source_env,
source_bucket,
ReplicationTargetOptions {
endpoint: &target.address(),
access_key: FAKE_ACCESS_KEY,
secret_key: FAKE_SECRET_KEY,
target_bucket: &target_bucket,
secure: false,
skip_tls_verify: false,
ca_cert_pem: None,
},
)
.await?;
put_bucket_replication_with_delete_statuses(source_env, source_bucket, &target_arn, "Enabled", Some("Enabled")).await?;
let target_client = fake_source_client(&target);
// Tag update on an existing version.
let tag_key = "ledger/tags.bin";
let tagged = source_client
.put_object()
.bucket(source_bucket)
.key(tag_key)
.body(ByteStream::from(payload(4 * 1024, 0x01)))
.send()
.await?;
let tag_source_version = tagged.version_id().ok_or("source PUT returned no version id")?.to_string();
assert_eq!(
wait_for_terminal_replication_status(&source_client, source_bucket, tag_key).await?,
"COMPLETED"
);
let tag_target_version = single_target_version(&target, &target_bucket, tag_key)?;
source_client
.put_object_tagging()
.bucket(source_bucket)
.key(tag_key)
.version_id(&tag_source_version)
.tagging(
Tagging::builder()
.tag_set(Tag::builder().key("phase").value("after").build()?)
.build()?,
)
.send()
.await?;
wait_until("tag update on the existing target version", || async {
let tags = target_client
.get_object_tagging()
.bucket(&target_bucket)
.key(tag_key)
.version_id(&tag_target_version)
.send()
.await?;
Ok(tags
.tag_set()
.iter()
.any(|tag| tag.key() == "phase" && tag.value() == "after"))
})
.await?;
assert_stable_single_version(&target, &target_bucket, tag_key, &tag_target_version).await?;
// Retention extension and legal hold on an existing version.
let lock_key = "ledger/lock.bin";
let locked = source_client
.put_object()
.bucket(source_bucket)
.key(lock_key)
.body(ByteStream::from(payload(4 * 1024, 0x02)))
.object_lock_mode(ObjectLockMode::Governance)
.object_lock_retain_until_date(retain_until())
.send()
.await?;
let lock_source_version = locked.version_id().ok_or("source PUT returned no version id")?.to_string();
assert_eq!(
wait_for_terminal_replication_status(&source_client, source_bucket, lock_key).await?,
"COMPLETED"
);
let lock_target_version = single_target_version(&target, &target_bucket, lock_key)?;
let extended = DateTime::from_secs(retain_until().secs() + 86_400);
source_client
.put_object_retention()
.bucket(source_bucket)
.key(lock_key)
.version_id(&lock_source_version)
.retention(
ObjectLockRetention::builder()
.mode(ObjectLockRetentionMode::Governance)
.retain_until_date(extended)
.build(),
)
.send()
.await?;
source_client
.put_object_legal_hold()
.bucket(source_bucket)
.key(lock_key)
.version_id(&lock_source_version)
.legal_hold(ObjectLockLegalHold::builder().status(ObjectLockLegalHoldStatus::On).build())
.send()
.await?;
wait_until("retention extension and legal hold on the existing target version", || async {
let head = target_client
.head_object()
.bucket(&target_bucket)
.key(lock_key)
.version_id(&lock_target_version)
.send()
.await?;
Ok(head.object_lock_retain_until_date().map(|date| date.secs()) == Some(extended.secs())
&& head.object_lock_legal_hold_status() == Some(&ObjectLockLegalHoldStatus::On))
})
.await?;
source_client
.put_object_legal_hold()
.bucket(source_bucket)
.key(lock_key)
.version_id(&lock_source_version)
.legal_hold(ObjectLockLegalHold::builder().status(ObjectLockLegalHoldStatus::Off).build())
.send()
.await?;
wait_until("legal hold removal on the existing target version", || async {
let head = target_client
.head_object()
.bucket(&target_bucket)
.key(lock_key)
.version_id(&lock_target_version)
.send()
.await?;
Ok(head.object_lock_legal_hold_status() == Some(&ObjectLockLegalHoldStatus::Off))
})
.await?;
assert_stable_single_version(&target, &target_bucket, lock_key, &lock_target_version).await?;
// Permanent delete of the older of two same-content generations.
let generations_key = "ledger/generations.bin";
let body = payload(4 * 1024, 0x03);
let older = source_client
.put_object()
.bucket(source_bucket)
.key(generations_key)
.body(ByteStream::from(body.clone()))
.send()
.await?;
let older_version = older.version_id().ok_or("source PUT returned no version id")?.to_string();
assert_eq!(
wait_for_terminal_replication_status(&source_client, source_bucket, generations_key).await?,
"COMPLETED"
);
let older_replica = single_target_version(&target, &target_bucket, generations_key)?;
source_client
.put_object()
.bucket(source_bucket)
.key(generations_key)
.body(ByteStream::from(body))
.send()
.await?;
assert_eq!(
wait_for_terminal_replication_status(&source_client, source_bucket, generations_key).await?,
"COMPLETED"
);
wait_until("both generations replicated", || async {
Ok(target.stored_versions(&target_bucket, generations_key).len() == 2)
})
.await?;
let newer_replica = target
.stored_versions(&target_bucket, generations_key)
.into_iter()
.map(|(version_id, _)| version_id)
.find(|version_id| version_id != &older_replica)
.ok_or("the second generation must have its own target version")?;
source_client
.delete_object()
.bucket(source_bucket)
.key(generations_key)
.version_id(&older_version)
.send()
.await?;
wait_until("permanent delete of the older generation's replica", || async {
let versions: Vec<String> = target
.stored_versions(&target_bucket, generations_key)
.into_iter()
.map(|(version_id, _)| version_id)
.collect();
Ok(versions == [newer_replica.clone()])
})
.await?;
assert_stable_single_version(&target, &target_bucket, generations_key, &newer_replica).await?;
// No mutation above may have gone out as a re-PUT: one upload per key.
for key in [tag_key, lock_key] {
let puts = target
.requests()
.iter()
.filter(|record| record.key.as_deref() == Some(key) && record.operation == FakeTargetOperation::PutObject)
.count();
assert_eq!(
puts, 1,
"{key}: a metadata update must not re-PUT the object on a target that mints its own ids"
);
}
target.shutdown().await;
Ok(())
}
fn single_target_version(target: &FakeS3Target, target_bucket: &str, key: &str) -> Result<String, Box<dyn Error + Send + Sync>> {
let versions = target.stored_versions(target_bucket, key);
match versions.as_slice() {
[(version_id, false)] => Ok(version_id.clone()),
other => Err(format!("{key}: expected exactly one live target version, got {other:?}").into()),
}
}
/// The target keeps holding exactly `version_id` for a few scanner cycles: a
/// re-driven PUT or a wrong delete would show up here.
async fn assert_stable_single_version(target: &FakeS3Target, target_bucket: &str, key: &str, version_id: &str) -> TestResult {
for _ in 0..8 {
let versions = target.stored_versions(target_bucket, key);
if versions.len() != 1 || versions[0].0 != version_id {
return Err(
format!("{key}: target versions drifted from the single expected replica {version_id}: {versions:?}").into(),
);
}
sleep(Duration::from_millis(500)).await;
}
Ok(())
}
async fn wait_until<F, Fut>(what: &str, mut probe: F) -> TestResult
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = Result<bool, Box<dyn Error + Send + Sync>>>,
{
let wait = async {
loop {
if probe().await? {
return Ok::<_, Box<dyn Error + Send + Sync>>(());
}
sleep(Duration::from_millis(250)).await;
}
};
timeout(Duration::from_secs(90), wait)
.await
.map_err(|_| format!("{what} did not happen within 90 seconds"))?
}
/// Wait until `key` is COMPLETED on the source and, for the observation
/// window after that, the target still holds exactly one live version of it.
async fn wait_for_replication_status_and_single_version(