fix(tier): reconcile paginated remote versions (#5405)

* feat(tiering): model provider version capabilities

* feat(tiering): persist opaque remote versions

* fix(tiering): use exact GCS generations

* fix(tiering): gate remote version state safely

* fix(tiering): gate remote version state writes

* fix(tiering): preserve remote version state on delete

* fix(tiering): accept unversioned transition responses

* fix(tiering): replay exact cleanup journals

* test(tiering): pin empty exact cleanup guard

* test(tiering): accept strict missing journal errors

* test(tiering): exercise free-version identity guard

* test(tiering): reach destination identity guard

* test(tiering): persist version identity drift

* fix(tier): reconcile paginated remote versions

* style(tier): format candidate validation test

* test(tiering): bind version drift fixture

---------

Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
cxymds
2026-07-29 23:04:36 +08:00
committed by GitHub
parent 09157485aa
commit ae11bcf2be
2 changed files with 158 additions and 69 deletions
@@ -172,7 +172,7 @@ impl ProviderVersionCapabilities {
} }
} }
fn validate_remote_version_id(version_id: &str) -> Result<(), Error> { pub(crate) fn validate_remote_version_id(version_id: &str) -> Result<(), Error> {
if version_id.is_empty() { if version_id.is_empty() {
return Err(Error::new( return Err(Error::new(
ErrorKind::InvalidData, ErrorKind::InvalidData,
@@ -29,6 +29,7 @@ use crate::client::{
api_remove::{RemoveObjectOptions, RemoveObjectResult}, api_remove::{RemoveObjectOptions, RemoveObjectResult},
api_s3_datatypes::ListVersionsResult, api_s3_datatypes::ListVersionsResult,
credentials::{Credentials, SignatureType, Static, Value}, credentials::{Credentials, SignatureType, Static, Value},
provider_versions::validate_remote_version_id,
transition_api::{BucketLookupType, Options, TransitionClient, TransitionCore}, transition_api::{BucketLookupType, Options, TransitionClient, TransitionCore},
transition_api::{ReadCloser, ReaderImpl}, transition_api::{ReadCloser, ReaderImpl},
}; };
@@ -189,44 +190,102 @@ impl WarmBackendS3 {
remote_bucket_versioning_from_status(config.status.as_ref().map(|status| status.as_str())) remote_bucket_versioning_from_status(config.status.as_ref().map(|status| status.as_str()))
} }
async fn list_transition_candidate_versions(&self, object: &str) -> Result<ListVersionsResult, std::io::Error> { async fn probe_transition_candidate_versions(
&self,
object: &str,
bucket_versioning: RemoteBucketVersioning,
) -> Result<TransitionCandidateProbe, std::io::Error> {
let remote_object = self.get_dest(object);
let mut opts = ListObjectsOptions::default(); let mut opts = ListObjectsOptions::default();
opts.set("prefix", &self.get_dest(object)); opts.set("prefix", &remote_object);
opts.set("max-keys", "2"); opts.set("max-keys", "1000");
self.client.list_object_versions_query(&self.bucket, &opts, "", "", "").await
let mut key_marker = String::new();
let mut version_id_marker = String::new();
let mut candidates = TransitionCandidateVersions::default();
loop {
let versions = self
.client
.list_object_versions_query(&self.bucket, &opts, &key_marker, &version_id_marker, "")
.await?;
candidates.extend(&remote_object, &versions);
if candidates.is_ambiguous() {
return Ok(TransitionCandidateProbe::Ambiguous);
}
if !versions.is_truncated {
return classify_transition_candidates(candidates, bucket_versioning);
}
advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?;
}
} }
} }
fn classify_transition_candidate_versions( fn classify_transition_candidates(
remote_object: &str, candidates: TransitionCandidateVersions,
bucket_versioning: RemoteBucketVersioning, bucket_versioning: RemoteBucketVersioning,
) -> Result<TransitionCandidateProbe, std::io::Error> {
let probe = candidates.classify(bucket_versioning);
if let TransitionCandidateProbe::VersionedPresent(version_id) = &probe {
validate_remote_version_id(version_id)?;
}
Ok(probe)
}
fn advance_version_markers(
key_marker: &mut String,
version_id_marker: &mut String,
versions: &ListVersionsResult, versions: &ListVersionsResult,
) -> TransitionCandidateProbe { ) -> Result<(), std::io::Error> {
if versions.is_truncated { let next_markers = (&versions.next_key_marker, &versions.next_version_id_marker);
return TransitionCandidateProbe::Ambiguous; if next_markers == (&*key_marker, &*version_id_marker) {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"ListObjectVersions pagination markers did not advance",
));
} }
key_marker.clone_from(&versions.next_key_marker);
version_id_marker.clone_from(&versions.next_version_id_marker);
Ok(())
}
if versions.delete_markers.iter().any(|marker| marker.key == remote_object) { #[derive(Default)]
return TransitionCandidateProbe::Ambiguous; struct TransitionCandidateVersions {
} version_id: Option<String>,
ambiguous: bool,
}
let mut exact_versions = versions.versions.iter().filter(|version| version.key == remote_object); impl TransitionCandidateVersions {
let Some(version) = exact_versions.next() else { fn extend(&mut self, remote_object: &str, versions: &ListVersionsResult) {
return TransitionCandidateProbe::Missing; for version in versions.versions.iter().filter(|version| version.key == remote_object) {
}; if self.version_id.is_some() {
if exact_versions.next().is_some() { self.ambiguous = true;
return TransitionCandidateProbe::Ambiguous; return;
} }
self.version_id = Some(version.version_id.clone());
match bucket_versioning {
RemoteBucketVersioning::Disabled => TransitionCandidateProbe::UnversionedPresent,
RemoteBucketVersioning::Suspended if version.version_id == "null" => {
TransitionCandidateProbe::VersionedPresent(version.version_id.clone())
} }
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled if !version.version_id.is_empty() => { }
TransitionCandidateProbe::VersionedPresent(version.version_id.clone())
fn is_ambiguous(&self) -> bool {
self.ambiguous
}
fn classify(self, bucket_versioning: RemoteBucketVersioning) -> TransitionCandidateProbe {
if self.ambiguous {
return TransitionCandidateProbe::Ambiguous;
}
let Some(version_id) = self.version_id else {
return TransitionCandidateProbe::Missing;
};
match bucket_versioning {
RemoteBucketVersioning::Disabled => TransitionCandidateProbe::UnversionedPresent,
RemoteBucketVersioning::Suspended if version_id == "null" => TransitionCandidateProbe::VersionedPresent(version_id),
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled if !version_id.is_empty() => {
TransitionCandidateProbe::VersionedPresent(version_id)
}
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled => TransitionCandidateProbe::Ambiguous,
} }
RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled => TransitionCandidateProbe::Ambiguous,
} }
} }
@@ -275,74 +334,109 @@ mod tests {
} }
} }
fn classify_pages(bucket_versioning: RemoteBucketVersioning, pages: &[ListVersionsResult]) -> TransitionCandidateProbe {
let mut candidates = TransitionCandidateVersions::default();
for page in pages {
candidates.extend("archive/object", page);
}
candidates.classify(bucket_versioning)
}
#[test] #[test]
fn transition_candidate_probe_classifier_is_fail_closed() { fn transition_candidate_probe_classifier_is_fail_closed() {
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(RemoteBucketVersioning::Disabled, &[list_versions(&[], &[], false)],),
"archive/object",
RemoteBucketVersioning::Disabled,
&list_versions(&[], &[], false),
),
TransitionCandidateProbe::Missing TransitionCandidateProbe::Missing
); );
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(RemoteBucketVersioning::Disabled, &[list_versions(&[("archive/object", "")], &[], false)],),
"archive/object",
RemoteBucketVersioning::Disabled,
&list_versions(&[("archive/object", "")], &[], false),
),
TransitionCandidateProbe::UnversionedPresent TransitionCandidateProbe::UnversionedPresent
); );
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(
"archive/object",
RemoteBucketVersioning::Enabled, RemoteBucketVersioning::Enabled,
&list_versions(&[("archive/object", "version-a")], &[], false), &[list_versions(&[("archive/object", "version-a")], &[], false)],
), ),
TransitionCandidateProbe::VersionedPresent("version-a".to_string()) TransitionCandidateProbe::VersionedPresent("version-a".to_string())
); );
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(
"archive/object",
RemoteBucketVersioning::Suspended, RemoteBucketVersioning::Suspended,
&list_versions(&[("archive/object", "null")], &[], false), &[list_versions(&[("archive/object", "null")], &[], false)],
), ),
TransitionCandidateProbe::VersionedPresent("null".to_string()) TransitionCandidateProbe::VersionedPresent("null".to_string())
); );
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(RemoteBucketVersioning::Enabled, &[list_versions(&[("archive/object", "")], &[], false)],),
"archive/object",
RemoteBucketVersioning::Enabled,
&list_versions(&[("archive/object", "")], &[], false),
),
TransitionCandidateProbe::Ambiguous TransitionCandidateProbe::Ambiguous
); );
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(
"archive/object",
RemoteBucketVersioning::Enabled, RemoteBucketVersioning::Enabled,
&list_versions(&[("archive/object", "version-a"), ("archive/object", "version-b")], &[], false), &[list_versions(
&[("archive/object", "version-a"), ("archive/object", "version-b")],
&[],
false,
)],
), ),
TransitionCandidateProbe::Ambiguous TransitionCandidateProbe::Ambiguous
); );
}
#[test]
fn transition_candidate_probe_reconciles_all_pages_and_ignores_delete_markers() {
assert_eq!(
classify_pages(
RemoteBucketVersioning::Enabled,
&[
list_versions(&[], &[("archive/object", "marker-a")], true),
list_versions(&[("archive/object", "version-a"), ("archive/object-adjacent", "unrelated"),], &[], false,),
],
),
TransitionCandidateProbe::VersionedPresent("version-a".to_string())
);
assert_eq!( assert_eq!(
classify_transition_candidate_versions( classify_pages(
"archive/object",
RemoteBucketVersioning::Enabled, RemoteBucketVersioning::Enabled,
&list_versions(&[("archive/object", "version-a")], &[("archive/object", "marker-a")], false), &[
), list_versions(&[("archive/object", "version-a")], &[], true),
TransitionCandidateProbe::Ambiguous list_versions(&[("archive/object", "version-b")], &[], false),
); ],
assert_eq!(
classify_transition_candidate_versions(
"archive/object",
RemoteBucketVersioning::Enabled,
&list_versions(&[("archive/object", "version-a")], &[], true),
), ),
TransitionCandidateProbe::Ambiguous TransitionCandidateProbe::Ambiguous
); );
} }
#[test]
fn transition_candidate_pagination_advances_both_markers() {
let mut key_marker = "old-key".to_string();
let mut version_id_marker = "old-version".to_string();
let page = ListVersionsResult {
next_key_marker: "next-key".to_string(),
next_version_id_marker: "next-version".to_string(),
..Default::default()
};
advance_version_markers(&mut key_marker, &mut version_id_marker, &page)
.expect("new ListObjectVersions markers should advance pagination");
assert_eq!(key_marker, "next-key");
assert_eq!(version_id_marker, "next-version");
let err = advance_version_markers(&mut key_marker, &mut version_id_marker, &page)
.expect_err("repeated ListObjectVersions markers must fail closed");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
}
#[test]
fn transition_candidate_probe_rejects_untrusted_version_ids() {
let mut candidates = TransitionCandidateVersions::default();
candidates.extend("archive/object", &list_versions(&[("archive/object", "version\ninjection")], &[], false));
let err = classify_transition_candidates(candidates, RemoteBucketVersioning::Enabled)
.expect_err("control characters in listed version IDs must fail closed");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
}
#[test] #[test]
fn remote_bucket_versioning_status_parser_fails_closed() { fn remote_bucket_versioning_status_parser_fails_closed() {
assert_eq!( assert_eq!(
@@ -397,12 +491,7 @@ impl WarmBackend for WarmBackendS3 {
async fn probe_transition_candidate(&self, object: &str) -> Result<TransitionCandidateProbe, std::io::Error> { async fn probe_transition_candidate(&self, object: &str) -> Result<TransitionCandidateProbe, std::io::Error> {
let bucket_versioning = self.remote_bucket_versioning().await?; let bucket_versioning = self.remote_bucket_versioning().await?;
let versions = self.list_transition_candidate_versions(object).await?; self.probe_transition_candidate_versions(object, bucket_versioning).await
Ok(classify_transition_candidate_versions(
&self.get_dest(object),
bucket_versioning,
&versions,
))
} }
async fn in_use(&self) -> Result<bool, std::io::Error> { async fn in_use(&self) -> Result<bool, std::io::Error> {