Compare commits

..

2 Commits

Author SHA1 Message Date
overtrue 782e80000d fix(ci): serialize performance on shared functional VMs 2026-09-05 19:34:52 +08:00
Zhengchao An f053862aad docs: request concrete behavior evidence in pull requests (#7196) 2026-09-05 18:56:26 +08:00
11 changed files with 207 additions and 942 deletions
+5 -5
View File
@@ -10,16 +10,16 @@ Use N/A when there is no related issue.
## Summary of Changes
<!--
Briefly explain what changed and why reviewers should accept it.
Focus on behavior, compatibility, and review-relevant context.
Describe the concrete problem and resulting behavior. For a behavior change, name the input or state that triggers it and the expected outcome. Explain any new dependency or abstraction that the change needs.
-->
## Verification
<!--
List the commands or checks you ran, for example:
- `make pre-commit`
Give 13 concrete pieces of evidence for the changed behavior: the test or command, its observed result, and the regression it catches. For a bug fix, record a failing-before/passing-after check or explain why it was unavailable.
Use N/A only when verification is not applicable.
Identify the tested commit and any local changes. When testing a prebuilt binary or external service, include its source/version and artifact identity; a successful run against a different build is not evidence for this change.
List relevant checks not run and the remaining risk. Use the validation tier in AGENTS.md; do not run broader checks solely to fill this section. For documentation-only changes, list the applicable documentation checks. Use N/A only when verification is not applicable.
-->
## Impact
+2 -15
View File
@@ -14,8 +14,8 @@
# Functional chain driver: runs the ten functional suites in a fixed order
# (upgrade -> s3 -> kms -> tier -> storage -> heal -> pool -> security ->
# replication, with performance on its own runner in parallel) and guarantees
# the chain keeps moving even when individual suites fail.
# replication -> performance). Each suite attempts the next handoff even
# when its tests fail.
#
# Each suite workflow can still be dispatched standalone (workflow_dispatch);
# only chain-triggered runs forward to the next suite via repository_dispatch,
@@ -59,16 +59,3 @@ jobs:
gh api --method POST repos/rustfs/rustfs/dispatches \
-f event_type='rustfs-chain-upgrade' \
-F 'client_payload[from_suite]=nightly-build'
- name: Dispatch performance suite (parallel, own runner)
env:
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
run: |
set -euo pipefail
if [ -z "${GH_TOKEN:-}" ]; then
echo "PF_TESTING_GH_TOKEN is not configured; cannot dispatch performance" >&2
exit 1
fi
gh api --method POST repos/rustfs/rustfs/dispatches \
-f event_type='rustfs-chain-performance' \
-F 'client_payload[from_suite]=nightly-build'
@@ -49,17 +49,16 @@ on:
type: boolean
default: true
repository_dispatch:
# Chain entry: dispatched by rustfs-functional-chain.yml (runs on its own
# pf-testing runner, in parallel with the shared-VM chain).
# Chain handoff: dispatched when the replication suite finishes.
types: [rustfs-chain-performance]
permissions:
contents: read
# Dedicated pf-testing runner/environment: own concurrency group so perf runs
# never block (or are blocked by) the pool-expansion / heal tests.
# The default performance nodes overlap the other suites' remote VMs, even
# though the runner differs. Hold the shared lock through cleanup as well.
concurrency:
group: rustfs-performance-test
group: rustfs-shared-functional-tests
cancel-in-progress: false
defaults:
+43 -10
View File
@@ -34,8 +34,7 @@ on:
- site
default: all
repository_dispatch:
# Chain handoff: dispatched when the security suite finishes. This is the
# last link of the functional chain.
# Chain handoff: dispatched when the security suite finishes.
types: [rustfs-chain-replication]
permissions:
@@ -62,9 +61,6 @@ env:
jobs:
replication-test:
runs-on: smoke-testing
# A failed replication run must not break the chain or the workflow: the
# failure is reported to rustfs/backlog instead (see the issue step).
continue-on-error: true
timeout-minutes: 360
if: ${{ github.event_name == 'workflow_dispatch' || github.event_name == 'repository_dispatch' }}
steps:
@@ -349,13 +345,50 @@ jobs:
'
done
- name: Chain complete
# Replication is the last link of the functional chain: nothing to
# dispatch after it. This step just records that the chain finished.
- name: "Continue functional chain (next: Performance)"
if: ${{ always() && github.event_name == 'repository_dispatch' }}
env:
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
run: |
echo "Functional chain complete: replication (final suite) finished."
echo "from_suite=security trigger=${{ github.event_name }} outcome=${{ steps.test.outcome }}"
set -uo pipefail
if [ -z "${GH_TOKEN:-}" ]; then
echo "PF_TESTING_GH_TOKEN is not configured; cannot dispatch the next suite" >&2
exit 1
fi
DISPATCHED=0
for attempt in 1 2 3; do
if gh api --method POST repos/rustfs/rustfs/dispatches \
-f event_type='rustfs-chain-performance' \
-F 'client_payload[from_suite]=replication'; then
echo "dispatched next suite Performance (attempt ${attempt})"
DISPATCHED=1
break
fi
echo "dispatch attempt ${attempt} failed; retrying in ${attempt}0s" >&2
sleep "${attempt}0"
done
if [ "${DISPATCHED:-0}" -ne 1 ]; then
echo "ERROR: functional chain stalled: could not dispatch Performance after 3 attempts" >&2
TITLE="[functional][chain] stalled after replication (run ${GITHUB_RUN_ID})"
BODY_FILE="$(mktemp)"
trap 'rm -f "${BODY_FILE}"' EXIT
{
echo "The functional chain could not hand off from **replication** to **Performance** after 3 attempts."
echo ""
echo "- Failed suite job: ${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}"
echo "- Expected next event: 'rustfs-chain-performance'"
echo "- Likely cause: PF_TESTING_GH_TOKEN lacks contents:write on rustfs/rustfs, or the GitHub API was unavailable."
echo "- Recovery: re-dispatch manually with"
FENCE="$(printf "\x60\x60\x60")"; echo " ${FENCE}"
echo " gh api --method POST repos/rustfs/rustfs/dispatches -f event_type='rustfs-chain-performance'"
FENCE="$(printf "\x60\x60\x60")"; echo " ${FENCE}"
} > "${BODY_FILE}"
gh issue create -R rustfs/backlog --title "${TITLE}" \
--body-file "${BODY_FILE}" --label functional-test \
|| gh issue create -R rustfs/backlog --title "${TITLE}" --body-file "${BODY_FILE}" \
|| echo "could not file the stall alert issue either; check the token" >&2
exit 1
fi
- name: Notify on failure
if: failure()
+2 -3
View File
@@ -89,9 +89,8 @@ pub mod bucket {
#[allow(clippy::module_inception)]
pub mod lifecycle {
pub use crate::bucket::lifecycle::lifecycle::{
Event, ExpirationOptions, IlmAction, LIFECYCLE_MALFORMED_XML_ERROR_KIND, Lifecycle, LifecycleCalculate,
ObjectOpts, RuleValidate, TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions, expected_expiry_time,
object_opts_from_object_info,
Event, ExpirationOptions, IlmAction, Lifecycle, LifecycleCalculate, ObjectOpts, RuleValidate,
TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions, expected_expiry_time, object_opts_from_object_info,
};
}
+3 -3
View File
@@ -15,9 +15,9 @@
use crate::object_api::ObjectInfo;
pub use rustfs_lifecycle::{
Event, ExpirationOptions, IlmAction, LIFECYCLE_MALFORMED_XML_ERROR_KIND, Lifecycle, LifecycleCalculate, ObjectOpts,
RuleValidate, TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions, abort_incomplete_multipart_upload_due,
expected_expiry_time, expiration_action_has_valid_target,
Event, ExpirationOptions, IlmAction, Lifecycle, LifecycleCalculate, ObjectOpts, RuleValidate, TRANSITION_COMPLETE,
TRANSITION_PENDING, TransitionOptions, abort_incomplete_multipart_upload_due, expected_expiry_time,
expiration_action_has_valid_target,
};
pub fn object_opts_from_object_info(oi: &ObjectInfo) -> ObjectOpts {
+30 -792
View File
@@ -18,7 +18,7 @@ use s3s::dto::{
BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, LifecycleRuleFilter,
NoncurrentVersionTransition, ObjectLockConfiguration, ObjectLockEnabled, RestoreRequest, Transition,
};
use std::collections::{HashMap, HashSet};
use std::collections::HashMap;
use std::sync::Arc;
use time::macros::offset;
use time::{self, Duration, OffsetDateTime};
@@ -61,66 +61,6 @@ const ERR_LIFECYCLE_EXPIRED_OBJECT_DELETE_MARKER_WITH_TAGS: &str =
"Rule with ExpiredObjectDeleteMarker cannot have tags based filtering";
const ERR_LIFECYCLE_RULE_MUST_HAVE_ACTION: &str = "Rule must have at least one of Expiration, Transition, NoncurrentVersionExpiration, NoncurrentVersionTransition, or DelMarkerExpiration";
const ERR_LIFECYCLE_PREFIX_FILTER_CONFLICT: &str = "Legacy Prefix and Filter cannot both be present in a lifecycle rule. Use Filter.Prefix instead of the top-level Prefix element.";
const ERR_LIFECYCLE_INVALID_NEWER_NONCURRENT_VERSIONS: &str = "'NewerNoncurrentVersions' must be a non-negative integer";
const ERR_LIFECYCLE_FILTER_TOO_MANY_PREDICATES: &str =
"Filter must have at most one of Prefix, Tag, ObjectSizeGreaterThan, ObjectSizeLessThan or And; combine predicates with And";
const ERR_LIFECYCLE_FILTER_AND_TOO_FEW_PREDICATES: &str = "Filter And must contain at least two predicates";
const ERR_LIFECYCLE_FILTER_DUPLICATE_TAG_KEY: &str = "Filter must not repeat a tag key";
const ERR_LIFECYCLE_FILTER_INVALID_TAG: &str = "Tag key must be 1-128 characters and tag value must be at most 256 characters";
const ERR_LIFECYCLE_FILTER_NEGATIVE_SIZE: &str = "ObjectSizeGreaterThan and ObjectSizeLessThan must not be negative";
const ERR_LIFECYCLE_FILTER_SIZE_RANGE: &str = "ObjectSizeGreaterThan must be smaller than ObjectSizeLessThan";
/// Longest tag key S3 accepts.
const MAX_TAG_KEY_LEN: usize = 128;
/// Longest tag value S3 accepts.
const MAX_TAG_VALUE_LEN: usize = 256;
/// A validation failure that the S3 boundary must answer with `MalformedXML`
/// rather than `InvalidArgument`: the document does not match the published
/// schema shape (wrong number of `Filter` predicates, a one-member `And`).
///
/// Everything else stays [`std::io::ErrorKind::Other`], which the boundary
/// already maps to `InvalidArgument`.
pub const LIFECYCLE_MALFORMED_XML_ERROR_KIND: std::io::ErrorKind = std::io::ErrorKind::InvalidData;
/// A persisted rule that could never have passed validation. Callers that can
/// report an error surface it; evaluation itself stays fail-closed and takes
/// no action for the rule.
pub const LIFECYCLE_CORRUPT_RULE_ERROR_KIND: std::io::ErrorKind = std::io::ErrorKind::InvalidData;
fn malformed_xml_error(message: &'static str) -> std::io::Error {
std::io::Error::new(LIFECYCLE_MALFORMED_XML_ERROR_KIND, message)
}
/// The retention count a rule keeps, or `None` when the persisted value is
/// negative — a shape PUT validation rejects, so reaching it means the rule
/// came from older persistence or an import.
///
/// A negative count must never be read as "retain everything": that is how an
/// invalid configuration silently stopped deleting versions (backlog#2201).
pub fn retained_noncurrent_versions(count: i32) -> Option<usize> {
usize::try_from(count).ok()
}
/// Does any rule carry a retention count that validation would have rejected?
pub fn lifecycle_has_corrupt_retention_count(lc: &BucketLifecycleConfiguration) -> bool {
lc.rules.iter().any(rule_has_corrupt_retention_count)
}
fn rule_has_corrupt_retention_count(rule: &LifecycleRule) -> bool {
let expiration_count = rule
.noncurrent_version_expiration
.as_ref()
.and_then(|expiration| expiration.newer_noncurrent_versions);
let transition_counts = rule
.noncurrent_version_transitions
.iter()
.flatten()
.filter_map(|transition| transition.newer_noncurrent_versions);
expiration_count
.into_iter()
.chain(transition_counts)
.any(|count| retained_noncurrent_versions(count).is_none())
}
pub use rustfs_scanner_metrics::metrics::IlmAction;
@@ -197,17 +137,6 @@ impl RuleValidate for LifecycleRule {
return Err(std::io::Error::other(ERR_LIFECYCLE_PREFIX_FILTER_CONFLICT));
}
if let Some(filter) = self.filter.as_ref() {
validate_lifecycle_filter(filter)?;
}
// A negative retention count was accepted and then read as "retain
// (almost) everything" during evaluation, so an HTTP-accepted rule
// silently stopped deleting versions (backlog#2201).
if rule_has_corrupt_retention_count(self) {
return Err(std::io::Error::other(ERR_LIFECYCLE_INVALID_NEWER_NONCURRENT_VERSIONS));
}
// Rule with DelMarkerExpiration cannot have tags based filtering
let has_tag_filter = self
.filter
@@ -240,14 +169,11 @@ impl RuleValidate for LifecycleRule {
// Rule must have at least one action
let has_expiration = self.expiration.is_some();
let has_transition = self.transitions.as_ref().is_some_and(|t| !t.is_empty());
// `NewerNoncurrentVersions` on its own is a MinIO extension, not an AWS
// form: it keeps the newest N noncurrent versions and expires the rest
// with no age condition. RustFS accepts it for MinIO compatibility, so
// it has to count as an action here — otherwise a count-only rule was
// rejected as actionless (backlog#2201).
let has_noncurrent_expiration = self.noncurrent_version_expiration.as_ref().is_some_and(|expiration| {
expiration.noncurrent_days.is_some() || expiration.newer_noncurrent_versions.is_some_and(|count| count > 0)
});
let has_noncurrent_expiration = self
.noncurrent_version_expiration
.as_ref()
.and_then(|e| e.noncurrent_days)
.is_some();
let has_noncurrent_transition = self
.noncurrent_version_transitions
.as_ref()
@@ -273,81 +199,6 @@ impl RuleValidate for LifecycleRule {
}
}
/// Structural validation for `LifecycleRuleFilter`.
///
/// The generated DTO is all-`Option`, so the S3 schema constraints have to be
/// checked here: at most one top-level predicate, an `And` that actually
/// combines at least two, no repeated tag key, tag key/value limits, and a
/// coherent non-negative size range (backlog#2201).
///
/// A filter with no predicate at all stays valid: AWS documents an empty
/// `Filter` as "applies to every object in the bucket", and rejecting it would
/// break the most common way to write an unconditional rule.
fn validate_lifecycle_filter(filter: &LifecycleRuleFilter) -> Result<(), std::io::Error> {
let top_level_predicates = usize::from(filter.prefix.is_some())
+ usize::from(filter.tag.is_some())
+ usize::from(filter.object_size_greater_than.is_some())
+ usize::from(filter.object_size_less_than.is_some())
+ usize::from(filter.and.is_some());
if top_level_predicates > 1 {
return Err(malformed_xml_error(ERR_LIFECYCLE_FILTER_TOO_MANY_PREDICATES));
}
if let Some(tag) = filter.tag.as_ref() {
validate_lifecycle_tag(tag)?;
}
if let Some(and) = filter.and.as_ref() {
let tags = and.tags.as_deref().unwrap_or(&[]);
let and_predicates = usize::from(and.prefix.is_some())
+ tags.len()
+ usize::from(and.object_size_greater_than.is_some())
+ usize::from(and.object_size_less_than.is_some());
if and_predicates < 2 {
return Err(malformed_xml_error(ERR_LIFECYCLE_FILTER_AND_TOO_FEW_PREDICATES));
}
let mut seen_keys = HashSet::with_capacity(tags.len());
for tag in tags {
validate_lifecycle_tag(tag)?;
let key = tag.key.as_deref().unwrap_or_default();
if !seen_keys.insert(key) {
return Err(std::io::Error::other(ERR_LIFECYCLE_FILTER_DUPLICATE_TAG_KEY));
}
}
validate_lifecycle_size_bounds(and.object_size_greater_than, and.object_size_less_than)?;
}
validate_lifecycle_size_bounds(filter.object_size_greater_than, filter.object_size_less_than)?;
Ok(())
}
/// S3 requires a tag to carry a key and value; both are length-bounded.
/// The DTO makes both optional, so incomplete tags have to be rejected here
/// rather than silently matching nothing.
fn validate_lifecycle_tag(tag: &s3s::dto::Tag) -> Result<(), std::io::Error> {
let key = tag.key.as_deref().unwrap_or_default();
let Some(value) = tag.value.as_deref() else {
return Err(std::io::Error::other(ERR_LIFECYCLE_FILTER_INVALID_TAG));
};
if key.is_empty() || key.chars().count() > MAX_TAG_KEY_LEN || value.chars().count() > MAX_TAG_VALUE_LEN {
return Err(std::io::Error::other(ERR_LIFECYCLE_FILTER_INVALID_TAG));
}
Ok(())
}
fn validate_lifecycle_size_bounds(greater_than: Option<i64>, less_than: Option<i64>) -> Result<(), std::io::Error> {
if greater_than.is_some_and(|size| size < 0) || less_than.is_some_and(|size| size < 0) {
return Err(std::io::Error::other(ERR_LIFECYCLE_FILTER_NEGATIVE_SIZE));
}
if let (Some(greater_than), Some(less_than)) = (greater_than, less_than)
&& greater_than >= less_than
{
return Err(std::io::Error::other(ERR_LIFECYCLE_FILTER_SIZE_RANGE));
}
Ok(())
}
fn lifecycle_rule_prefix(rule: &LifecycleRule) -> Option<&str> {
// Prefer a non-empty legacy prefix; treat an empty legacy prefix as if it were not set
if let Some(p) = rule.prefix.as_deref()
@@ -438,10 +289,6 @@ impl Lifecycle for BucketLifecycleConfiguration {
{
return true;
}
// A positive count is an action on its own (the MinIO count-only
// form). Zero means "no count constraint" here, exactly as the
// batch limit path reads it, and a negative count is corrupt —
// neither makes the rule active (backlog#2201).
if let Some(newer_noncurrent_versions) = rule_noncurrent_version_expiration.newer_noncurrent_versions
&& newer_noncurrent_versions > 0
{
@@ -710,23 +557,6 @@ impl Lifecycle for BucketLifecycleConfiguration {
if let Some(ref lc_rules) = self.filter_rules(obj).await {
for rule in lc_rules.iter() {
// A retention count that PUT validation would have rejected can
// only come from older persistence or an import. Take no action
// for the rule instead of allowing another action on the same
// corrupt rule to delete or transition an object (backlog#2201).
if rule_has_corrupt_retention_count(rule) {
debug!(
event = EVENT_LIFECYCLE_NONCURRENT_EXPIRY_SKIPPED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
object = %obj.name,
rule_id = %rule.id.clone().unwrap_or_default(),
reason = "corrupt_newer_noncurrent_versions",
"Skipped lifecycle evaluation for a rule with an invalid retention count"
);
continue;
}
if obj.is_latest && obj.expired_object_deletemarker() {
if let Some(expiration) = rule.expiration.as_ref()
&& expiration.expired_object_delete_marker.is_some_and(|v| v)
@@ -784,23 +614,15 @@ impl Lifecycle for BucketLifecycleConfiguration {
if !obj.is_latest
&& let Some(ref noncurrent_version_expiration) = rule.noncurrent_version_expiration
&& let Some(retain_newer_noncurrent_versions) = noncurrent_version_expiration.newer_noncurrent_versions
&& let Some(retained) = retained_noncurrent_versions(retain_newer_noncurrent_versions)
&& newer_noncurrent_versions < retained
&& newer_noncurrent_versions < usize::try_from(retain_newer_noncurrent_versions).unwrap_or(usize::MAX)
{
continue;
}
if !obj.is_latest
&& let Some(ref noncurrent_version_expiration) = rule.noncurrent_version_expiration
&& (noncurrent_version_expiration.noncurrent_days.is_some()
|| noncurrent_version_expiration
.newer_noncurrent_versions
.is_some_and(|count| count > 0))
&& let Some(noncurrent_days) = noncurrent_version_expiration.noncurrent_days
{
// A count-only rule (MinIO extension) has no age condition:
// every version past the retained count is due as soon as it
// became noncurrent, i.e. zero days after the successor.
let noncurrent_days = noncurrent_version_expiration.noncurrent_days.unwrap_or(0);
if let Some(successor_mod_time) = obj.successor_mod_time {
let expected_expiry = expected_expiry_time(successor_mod_time, noncurrent_days);
if now.unix_timestamp() >= expected_expiry.unix_timestamp() {
@@ -963,18 +785,15 @@ impl Lifecycle for BucketLifecycleConfiguration {
for rule in filter_rules.iter() {
if let Some(ref noncurrent_version_expiration) = rule.noncurrent_version_expiration {
return if let Some(newer_noncurrent_versions) = noncurrent_version_expiration.newer_noncurrent_versions {
// Zero means "no count constraint"; a negative count is
// corrupt and must not be read as "retain everything"
// (backlog#2201). Neither yields a limit event.
let Some(retained) = retained_noncurrent_versions(newer_noncurrent_versions).filter(|c| *c > 0) else {
if newer_noncurrent_versions == 0 {
continue;
};
}
Event {
action: IlmAction::DeleteVersionAction,
rule_id: rule.id.clone().unwrap_or_default(),
noncurrent_days: u32::try_from(noncurrent_version_expiration.noncurrent_days.unwrap_or(0))
.unwrap_or(u32::MAX),
newer_noncurrent_versions: retained,
newer_noncurrent_versions: usize::try_from(newer_noncurrent_versions).unwrap_or(usize::MAX),
due: Some(OffsetDateTime::UNIX_EPOCH),
storage_class: "".into(),
}
@@ -1319,11 +1138,7 @@ mod tests {
use super::*;
use metrics_util::MetricKind;
use metrics_util::debugging::{DebugValue, DebuggingRecorder};
use s3s::dto::{
LifecycleRuleAndOperator, LifecycleRuleFilter, NoncurrentVersionExpiration, NoncurrentVersionTransition,
TransitionStorageClass,
};
use s3s::xml::{Deserialize as XmlDeserialize, SerializeContent as XmlSerializeContent};
use s3s::dto::{LifecycleRuleFilter, TransitionStorageClass};
use serial_test::serial;
use std::sync::Arc;
use time::macros::datetime;
@@ -4344,601 +4159,24 @@ mod tests {
assert_eq!(event.action, IlmAction::NoneAction);
}
// Property-based tests for the rule evaluator (backlog#1148 ilm-14,
// follow-up to backlog#1030 / rustfs#4455).
//
// backlog#1030 found that the old hand-written winner comparator was not a
// strict weak ordering and could panic — proof that enumerated cases do
// not cover the space of colliding events. These properties pin, over
// randomized rule sets and object states:
//
// * `eval_inner` never panics and is deterministic for a fixed input;
// * the winning event matches an independently recomputed candidate set:
// earliest `due` wins, ties break toward delete-class actions (the
// `min_by_key` selection that replaced the rustfs#4455 comparator);
// * `expected_expiry_time` is monotonically non-decreasing in `days` and
// always lands on the processing boundary, both at production defaults
// and under an explicit `RUSTFS_ILM_PROCESS_TIME`.
//
// Case counts are tuned so the whole module runs in seconds inside the
// default CI test job.
// ---- backlog#2201: retention-count and Filter invariants -----------------
fn rule_with_noncurrent_expiration(expiration: NoncurrentVersionExpiration) -> LifecycleRule {
LifecycleRule {
status: ExpirationStatus::from_static(ExpirationStatus::ENABLED),
expiration: None,
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
filter: None,
id: Some("noncurrent".to_string()),
noncurrent_version_expiration: Some(expiration),
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}
}
fn rule_with_filter(filter: LifecycleRuleFilter) -> LifecycleRule {
LifecycleRule {
status: ExpirationStatus::from_static(ExpirationStatus::ENABLED),
expiration: Some(LifecycleExpiration {
days: Some(1),
..Default::default()
}),
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
filter: Some(filter),
id: Some("filtered".to_string()),
noncurrent_version_expiration: None,
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}
}
fn config_with_rules(rules: Vec<LifecycleRule>) -> BucketLifecycleConfiguration {
BucketLifecycleConfiguration {
expiry_updated_at: None,
rules,
}
}
fn tag(key: &str, value: &str) -> s3s::dto::Tag {
s3s::dto::Tag {
key: Some(key.to_string()),
value: Some(value.to_string()),
}
}
#[tokio::test]
async fn validate_rejects_negative_newer_noncurrent_versions() {
// A negative retention count used to be accepted and then read as
// usize::MAX during evaluation, so the rule silently stopped deleting
// versions (backlog#2201).
let lc = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(30),
newer_noncurrent_versions: Some(-1),
})]);
let err = lc
.validate(&ObjectLockConfiguration::default())
.await
.expect_err("a negative retention count must be rejected");
assert_eq!(err.to_string(), ERR_LIFECYCLE_INVALID_NEWER_NONCURRENT_VERSIONS);
assert_ne!(err.kind(), LIFECYCLE_MALFORMED_XML_ERROR_KIND, "value errors stay InvalidArgument");
}
#[tokio::test]
async fn validate_rejects_negative_newer_noncurrent_versions_on_transition() {
let mut rule = rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(30),
newer_noncurrent_versions: None,
});
rule.noncurrent_version_transitions = Some(vec![NoncurrentVersionTransition {
newer_noncurrent_versions: Some(-3),
noncurrent_days: Some(1),
storage_class: Some(TransitionStorageClass::from_static(TransitionStorageClass::GLACIER)),
}]);
// The transition validator already refuses a negative count, and it runs
// first, so this pins the rejection rather than the message. The gap
// this PR closes is the expiration side, which had no such check.
config_with_rules(vec![rule])
.validate(&ObjectLockConfiguration::default())
.await
.expect_err("a negative retention count on a transition must be rejected");
}
#[tokio::test]
async fn zero_newer_noncurrent_versions_means_no_count_constraint() {
// Zero carries no constraint, matching how the batch limit path has
// always read it. Alongside an age condition the rule is valid; on its
// own it says nothing, so the rule has no action.
config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(30),
newer_noncurrent_versions: Some(0),
})])
.validate(&ObjectLockConfiguration::default())
.await
.expect("zero count alongside NoncurrentDays is valid");
let err = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: None,
newer_noncurrent_versions: Some(0),
})])
.validate(&ObjectLockConfiguration::default())
.await
.expect_err("a zero count on its own is not an action");
assert_eq!(err.to_string(), ERR_LIFECYCLE_RULE_MUST_HAVE_ACTION);
}
#[tokio::test]
async fn validate_accepts_count_only_noncurrent_expiration() {
// MinIO extension: NewerNoncurrentVersions with no NoncurrentDays. It
// used to be rejected as an actionless rule (backlog#2201).
let lc = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: None,
newer_noncurrent_versions: Some(2),
})]);
lc.validate(&ObjectLockConfiguration::default())
.await
.expect("a count-only noncurrent expiration rule is accepted");
}
#[tokio::test]
async fn eval_inner_expires_versions_beyond_count_only_retention() {
// Count-only rules have no age condition: everything past the retained
// count is due as soon as it became noncurrent.
let lc = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: None,
newer_noncurrent_versions: Some(2),
})]);
let opts = ObjectOpts {
name: "obj".to_string(),
mod_time: Some(datetime!(2025-01-15 10:30:45 UTC)),
successor_mod_time: Some(datetime!(2025-01-15 10:30:45 UTC)),
is_latest: false,
num_versions: 5,
..Default::default()
};
// Rank 2 is the third-newest noncurrent version: past a retention of 2.
let expired = lc.eval_inner(&opts, datetime!(2025-01-15 10:30:46 UTC), 2).await;
assert_eq!(expired.action, IlmAction::DeleteVersionAction);
assert_eq!(expired.rule_id, "noncurrent");
// Rank 1 is still within the retained count.
let retained = lc.eval_inner(&opts, datetime!(2025-01-15 10:30:46 UTC), 1).await;
assert_eq!(retained.action, IlmAction::NoneAction);
}
#[tokio::test]
#[serial]
async fn eval_inner_keeps_age_condition_when_count_and_days_are_set() {
// With both set, the count gates which versions are candidates and the
// age condition still decides when they are due.
with_default_ilm_process_time(|| {});
let lc = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(10),
newer_noncurrent_versions: Some(1),
})]);
let opts = ObjectOpts {
name: "obj".to_string(),
mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
successor_mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
is_latest: false,
num_versions: 3,
..Default::default()
};
let too_young = lc.eval_inner(&opts, datetime!(2025-01-05 00:00:00 UTC), 2).await;
assert_eq!(too_young.action, IlmAction::NoneAction, "the age condition still applies");
let due = lc.eval_inner(&opts, datetime!(2025-01-20 00:00:00 UTC), 2).await;
assert_eq!(due.action, IlmAction::DeleteVersionAction);
}
#[tokio::test]
async fn eval_inner_takes_no_action_for_a_corrupt_retention_count() {
// Reachable only from older persistence or an import; it must not be
// read as "retain everything", and it must not delete either.
let lc = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(-1),
})]);
let opts = ObjectOpts {
name: "obj".to_string(),
mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
successor_mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
is_latest: false,
num_versions: 3,
..Default::default()
};
let event = lc.eval_inner(&opts, datetime!(2025-06-01 00:00:00 UTC), 2).await;
assert_eq!(event.action, IlmAction::NoneAction);
}
#[tokio::test]
async fn eval_inner_does_not_expire_latest_object_for_a_corrupt_retention_rule() {
let mut rule = rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(-1),
});
rule.expiration = Some(LifecycleExpiration {
days: Some(1),
..Default::default()
});
let lc = config_with_rules(vec![rule]);
let opts = ObjectOpts {
name: "obj".to_string(),
mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
is_latest: true,
..Default::default()
};
let event = lc.eval_inner(&opts, datetime!(2025-06-01 00:00:00 UTC), 0).await;
assert_eq!(event.action, IlmAction::NoneAction);
}
#[tokio::test]
async fn eval_inner_does_not_delete_latest_marker_for_a_corrupt_retention_rule() {
let mut expired_marker_rule = rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(-1),
});
expired_marker_rule.expiration = Some(LifecycleExpiration {
expired_object_delete_marker: Some(true),
..Default::default()
});
let mut aged_marker_rule = rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(-1),
});
aged_marker_rule.del_marker_expiration = Some(s3s::dto::DelMarkerExpiration { days: Some(1) });
for rule in [expired_marker_rule, aged_marker_rule] {
let lc = config_with_rules(vec![rule]);
let opts = ObjectOpts {
name: "obj".to_string(),
mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
version_id: Some(Uuid::new_v4()),
is_latest: true,
delete_marker: true,
num_versions: 1,
..Default::default()
};
let event = lc.eval_inner(&opts, datetime!(2025-06-01 00:00:00 UTC), 0).await;
assert_eq!(event.action, IlmAction::NoneAction);
}
}
#[test]
fn corrupt_retention_count_is_detected_on_either_action() {
let mut transition_rule = rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(0),
});
transition_rule.noncurrent_version_transitions = Some(vec![NoncurrentVersionTransition {
newer_noncurrent_versions: Some(-1),
noncurrent_days: Some(1),
storage_class: Some(TransitionStorageClass::from_static(TransitionStorageClass::GLACIER)),
}]);
assert!(lifecycle_has_corrupt_retention_count(&config_with_rules(vec![
rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(-1),
})
])));
assert!(lifecycle_has_corrupt_retention_count(&config_with_rules(vec![transition_rule])));
assert!(!lifecycle_has_corrupt_retention_count(&config_with_rules(vec![
rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(3),
})
])));
}
#[test]
fn count_only_rules_are_active_only_for_a_positive_count() {
let positive = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: None,
newer_noncurrent_versions: Some(2),
})]);
assert!(positive.has_active_rules(""));
let corrupt = config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: None,
newer_noncurrent_versions: Some(-1),
})]);
assert!(!corrupt.has_active_rules(""), "a corrupt retention count must not make a rule active");
}
#[tokio::test]
async fn noncurrent_versions_expiration_limit_ignores_a_corrupt_count() {
// The batch path must not read a negative count as "retain everything".
let lc = Arc::new(config_with_rules(vec![rule_with_noncurrent_expiration(NoncurrentVersionExpiration {
noncurrent_days: Some(1),
newer_noncurrent_versions: Some(-1),
})]));
let opts = ObjectOpts {
name: "obj".to_string(),
mod_time: Some(datetime!(2025-01-01 00:00:00 UTC)),
is_latest: false,
..Default::default()
};
let event = lc.noncurrent_versions_expiration_limit(&opts).await;
assert_eq!(event.action, IlmAction::NoneAction);
assert_eq!(event.newer_noncurrent_versions, 0);
}
#[tokio::test]
async fn validate_covers_filter_invariants() {
struct Case {
name: &'static str,
filter: LifecycleRuleFilter,
expected: Option<(&'static str, std::io::ErrorKind)>,
}
let cases = vec![
Case {
// AWS documents an empty Filter as "every object in the bucket".
name: "empty filter applies to all objects",
filter: LifecycleRuleFilter::default(),
expected: None,
},
Case {
name: "single prefix predicate",
filter: LifecycleRuleFilter {
prefix: Some("logs/".to_string()),
..Default::default()
},
expected: None,
},
Case {
name: "two top-level predicates",
filter: LifecycleRuleFilter {
prefix: Some("logs/".to_string()),
tag: Some(tag("env", "prod")),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_TOO_MANY_PREDICATES, LIFECYCLE_MALFORMED_XML_ERROR_KIND)),
},
Case {
name: "prefix alongside And",
filter: LifecycleRuleFilter {
prefix: Some("logs/".to_string()),
and: Some(LifecycleRuleAndOperator {
prefix: Some("logs/".to_string()),
tags: Some(vec![tag("env", "prod")]),
..Default::default()
}),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_TOO_MANY_PREDICATES, LIFECYCLE_MALFORMED_XML_ERROR_KIND)),
},
Case {
name: "And with a single member",
filter: LifecycleRuleFilter {
and: Some(LifecycleRuleAndOperator {
prefix: Some("logs/".to_string()),
..Default::default()
}),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_AND_TOO_FEW_PREDICATES, LIFECYCLE_MALFORMED_XML_ERROR_KIND)),
},
Case {
name: "And with two members",
filter: LifecycleRuleFilter {
and: Some(LifecycleRuleAndOperator {
prefix: Some("logs/".to_string()),
tags: Some(vec![tag("env", "prod")]),
..Default::default()
}),
..Default::default()
},
expected: None,
},
Case {
name: "And with two tags",
filter: LifecycleRuleFilter {
and: Some(LifecycleRuleAndOperator {
tags: Some(vec![tag("env", "prod"), tag("team", "storage")]),
..Default::default()
}),
..Default::default()
},
expected: None,
},
Case {
name: "And repeating a tag key",
filter: LifecycleRuleFilter {
and: Some(LifecycleRuleAndOperator {
tags: Some(vec![tag("env", "prod"), tag("env", "dev")]),
..Default::default()
}),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_DUPLICATE_TAG_KEY, std::io::ErrorKind::Other)),
},
Case {
name: "empty tag key",
filter: LifecycleRuleFilter {
tag: Some(tag("", "prod")),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_INVALID_TAG, std::io::ErrorKind::Other)),
},
Case {
name: "missing tag key",
filter: LifecycleRuleFilter {
tag: Some(s3s::dto::Tag {
key: None,
value: Some("prod".to_string()),
}),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_INVALID_TAG, std::io::ErrorKind::Other)),
},
Case {
name: "missing tag value",
filter: LifecycleRuleFilter {
tag: Some(s3s::dto::Tag {
key: Some("env".to_string()),
value: None,
}),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_INVALID_TAG, std::io::ErrorKind::Other)),
},
Case {
name: "empty tag value",
filter: LifecycleRuleFilter {
tag: Some(tag("env", "")),
..Default::default()
},
expected: None,
},
Case {
name: "tag key at the limit",
filter: LifecycleRuleFilter {
tag: Some(tag(&"k".repeat(MAX_TAG_KEY_LEN), "prod")),
..Default::default()
},
expected: None,
},
Case {
name: "tag key past the limit",
filter: LifecycleRuleFilter {
tag: Some(tag(&"k".repeat(MAX_TAG_KEY_LEN + 1), "prod")),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_INVALID_TAG, std::io::ErrorKind::Other)),
},
Case {
name: "tag value past the limit",
filter: LifecycleRuleFilter {
tag: Some(tag("env", &"v".repeat(MAX_TAG_VALUE_LEN + 1))),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_INVALID_TAG, std::io::ErrorKind::Other)),
},
Case {
name: "negative ObjectSizeGreaterThan",
filter: LifecycleRuleFilter {
object_size_greater_than: Some(-1),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_NEGATIVE_SIZE, std::io::ErrorKind::Other)),
},
Case {
name: "negative ObjectSizeLessThan",
filter: LifecycleRuleFilter {
object_size_less_than: Some(-5),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_NEGATIVE_SIZE, std::io::ErrorKind::Other)),
},
Case {
name: "inverted size range inside And",
filter: LifecycleRuleFilter {
and: Some(LifecycleRuleAndOperator {
object_size_greater_than: Some(100),
object_size_less_than: Some(100),
..Default::default()
}),
..Default::default()
},
expected: Some((ERR_LIFECYCLE_FILTER_SIZE_RANGE, std::io::ErrorKind::Other)),
},
Case {
name: "valid size range inside And",
filter: LifecycleRuleFilter {
and: Some(LifecycleRuleAndOperator {
object_size_greater_than: Some(1),
object_size_less_than: Some(2),
..Default::default()
}),
..Default::default()
},
expected: None,
},
];
for case in cases {
let result = config_with_rules(vec![rule_with_filter(case.filter)])
.validate(&ObjectLockConfiguration::default())
.await;
match (case.expected, result) {
(None, Ok(())) => {}
(None, Err(err)) => panic!("{}: expected acceptance, got {err}", case.name),
(Some((message, _)), Ok(())) => panic!("{}: expected rejection with {message}", case.name),
(Some((message, kind)), Err(err)) => {
assert_eq!(err.to_string(), message, "{}", case.name);
assert_eq!(err.kind(), kind, "{}: wrong S3 error category", case.name);
}
}
}
}
#[tokio::test]
async fn validate_keeps_legacy_prefix_and_filter_mutually_exclusive() {
let mut rule = rule_with_filter(LifecycleRuleFilter {
prefix: Some("logs/".to_string()),
..Default::default()
});
rule.prefix = Some("legacy/".to_string());
let err = config_with_rules(vec![rule])
.validate(&ObjectLockConfiguration::default())
.await
.expect_err("legacy Prefix and Filter cannot both be present");
assert_eq!(err.to_string(), ERR_LIFECYCLE_PREFIX_FILTER_CONFLICT);
}
#[test]
fn count_only_rule_round_trips_through_xml() {
// The MinIO count-only form has to survive the wire codec, or the rule
// this PR now accepts could not be persisted and read back.
let xml = br#"<LifecycleConfiguration><Rule><ID>count-only</ID><Status>Enabled</Status><Filter></Filter><NoncurrentVersionExpiration><NewerNoncurrentVersions>2</NewerNoncurrentVersions></NoncurrentVersionExpiration></Rule></LifecycleConfiguration>"#;
let mut deserializer = s3s::xml::Deserializer::new(xml);
let parsed =
<BucketLifecycleConfiguration as XmlDeserialize>::deserialize(&mut deserializer).expect("count-only XML parses");
let expiration = parsed.rules[0]
.noncurrent_version_expiration
.as_ref()
.expect("noncurrent expiration is present");
assert_eq!(expiration.newer_noncurrent_versions, Some(2));
assert_eq!(expiration.noncurrent_days, None);
let mut buf = Vec::new();
let mut serializer = s3s::xml::Serializer::new(&mut buf);
XmlSerializeContent::serialize_content(&parsed, &mut serializer).expect("count-only config serializes");
let serialized = String::from_utf8(buf).expect("serialized XML is UTF-8");
assert!(
serialized.contains("<NewerNoncurrentVersions>2</NewerNoncurrentVersions>"),
"retention count survives the round trip: {serialized}"
);
assert!(
!serialized.contains("<NoncurrentDays>"),
"a count-only rule must not gain an age condition: {serialized}"
);
}
/// Property-based tests for the rule evaluator (backlog#1148 ilm-14,
/// follow-up to backlog#1030 / rustfs#4455).
///
/// backlog#1030 found that the old hand-written winner comparator was not a
/// strict weak ordering and could panic — proof that enumerated cases do
/// not cover the space of colliding events. These properties pin, over
/// randomized rule sets and object states:
///
/// * `eval_inner` never panics and is deterministic for a fixed input;
/// * the winning event matches an independently recomputed candidate set:
/// earliest `due` wins, ties break toward delete-class actions (the
/// `min_by_key` selection that replaced the rustfs#4455 comparator);
/// * `expected_expiry_time` is monotonically non-decreasing in `days` and
/// always lands on the processing boundary, both at production defaults
/// and under an explicit `RUSTFS_ILM_PROCESS_TIME`.
///
/// Case counts are tuned so the whole module runs in seconds inside the
/// default CI test job.
mod proptests {
use super::*;
use proptest::prelude::*;
+1 -15
View File
@@ -22,10 +22,7 @@ use rustfs_replication::ReplicationStatusType;
use rustfs_scanner_metrics::metrics::IlmAction;
use crate::object_lock;
use crate::{
Event, LIFECYCLE_CORRUPT_RULE_ERROR_KIND, Lifecycle, ObjectOpts, expiration_action_has_valid_target,
lifecycle_has_corrupt_retention_count,
};
use crate::{Event, Lifecycle, ObjectOpts, expiration_action_has_valid_target};
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
const LOG_SUBSYSTEM_LIFECYCLE: &str = "lifecycle";
@@ -158,17 +155,6 @@ impl Evaluator {
format!("number of versions mismatch, expected {}, got {}", objs[0].num_versions, objs.len()),
));
}
// PUT validation rejects a negative retention count, so a rule that
// carries one came from older persistence or an import. Report it
// instead of evaluating a configuration that cannot be honoured;
// `eval_inner` independently takes no action for such a rule
// (backlog#2201).
if lifecycle_has_corrupt_retention_count(&self.policy) {
return Err(std::io::Error::new(
LIFECYCLE_CORRUPT_RULE_ERROR_KIND,
"lifecycle configuration carries a negative 'NewerNoncurrentVersions'",
));
}
Ok(self.eval_inner(objs, OffsetDateTime::now_utc()).await)
}
}
+3 -82
View File
@@ -27,8 +27,8 @@ use super::storage_api::bucket_usecase::bucket::target::BucketTarget;
use super::storage_api::bucket_usecase::bucket::{
ObjectLockConfigExt as _, VersioningConfigExt as _,
lifecycle::bucket_lifecycle_ops::{
LIFECYCLE_MALFORMED_XML_ERROR_KIND, enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects,
run_stale_multipart_upload_cleanup_once, validate_lifecycle_config, validate_transition_tier,
enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects, run_stale_multipart_upload_cleanup_once,
validate_lifecycle_config, validate_transition_tier,
},
metadata::{
BUCKET_CORS_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG,
@@ -1188,21 +1188,6 @@ fn validate_lifecycle_rule_status(rules: &[LifecycleRule]) -> std::result::Resul
Ok(())
}
/// Map a lifecycle validation failure onto the S3 error the client should see.
///
/// The validator reports a schema-shape violation (a `Filter` with more than
/// one predicate, a one-member `And`) with
/// [`LIFECYCLE_MALFORMED_XML_ERROR_KIND`]; AWS answers those with
/// `MalformedXML`. Everything else is a value the schema allows but S3 refuses,
/// which stays `InvalidArgument` — the code this path has always returned
/// (backlog#2201).
fn lifecycle_validation_error(err: &std::io::Error) -> S3Error {
if err.kind() == LIFECYCLE_MALFORMED_XML_ERROR_KIND {
return S3Error::with_message(S3ErrorCode::MalformedXML, format!("Malformed XML: {err}"));
}
s3_error!(InvalidArgument, "{err}")
}
fn lifecycle_has_transition_rules(config: &BucketLifecycleConfiguration) -> bool {
config.rules.iter().any(|rule| {
rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED)
@@ -2311,7 +2296,7 @@ impl DefaultBucketUsecase {
};
if let Err(err) = validate_lifecycle_config(&input_cfg, &rcfg).await {
return Err(lifecycle_validation_error(&err));
return Err(s3_error!(InvalidArgument, "{err}"));
}
if let Err(err) = validate_transition_tier(&input_cfg).await {
@@ -4026,70 +4011,6 @@ mod tests {
assert_eq!(rules[2].id.as_deref(), Some("rule-2"));
}
#[tokio::test]
async fn put_bucket_lifecycle_validation_errors_keep_their_s3_code() {
// The PUT path answers a schema-shape violation with MalformedXML and a
// rejected value with InvalidArgument. Both categories are produced by
// the real validator here, so the mapping cannot drift from it
// (backlog#2201).
let malformed = validate_lifecycle_config(
&BucketLifecycleConfiguration {
expiry_updated_at: None,
rules: vec![LifecycleRule {
status: ExpirationStatus::from_static(ExpirationStatus::ENABLED),
expiration: Some(LifecycleExpiration {
days: Some(1),
..Default::default()
}),
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
filter: Some(s3s::dto::LifecycleRuleFilter {
prefix: Some("logs/".to_string()),
tag: Some(s3s::dto::Tag {
key: Some("env".to_string()),
value: Some("prod".to_string()),
}),
..Default::default()
}),
id: Some("two-predicates".to_string()),
noncurrent_version_expiration: None,
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}],
},
&ObjectLockConfiguration::default(),
)
.await
.expect_err("a Filter with two predicates is a schema violation");
assert_eq!(*lifecycle_validation_error(&malformed).code(), S3ErrorCode::MalformedXML);
let invalid_value = validate_lifecycle_config(
&BucketLifecycleConfiguration {
expiry_updated_at: None,
rules: vec![LifecycleRule {
status: ExpirationStatus::from_static(ExpirationStatus::ENABLED),
expiration: None,
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
filter: None,
id: Some("negative-count".to_string()),
noncurrent_version_expiration: Some(s3s::dto::NoncurrentVersionExpiration {
noncurrent_days: Some(30),
newer_noncurrent_versions: Some(-1),
}),
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}],
},
&ObjectLockConfiguration::default(),
)
.await
.expect_err("a negative retention count is rejected");
assert_eq!(*lifecycle_validation_error(&invalid_value).code(), S3ErrorCode::InvalidArgument);
}
#[test]
fn validate_lifecycle_rule_status_rejects_invalid_status() {
let rules = vec![LifecycleRule {
-6
View File
@@ -392,12 +392,6 @@ pub(crate) mod bucket {
lc.validate(lock_config).await
}
/// The `std::io::ErrorKind` [`validate_lifecycle_config`] uses for a
/// lifecycle document that violates the published schema shape, which
/// the S3 boundary answers with `MalformedXML` (backlog#2201).
pub(crate) const LIFECYCLE_MALFORMED_XML_ERROR_KIND: std::io::ErrorKind =
crate::storage::storage_api::ecstore_bucket::lifecycle::lifecycle::LIFECYCLE_MALFORMED_XML_ERROR_KIND;
}
pub(crate) mod lifecycle_contract {
+114 -6
View File
@@ -1,5 +1,5 @@
#!/usr/bin/env python3
"""Run the security workflow's evidence and result steps without remote VMs."""
"""Exercise functional chain dispatch and security evidence without remote VMs."""
from __future__ import annotations
@@ -18,16 +18,20 @@ WORKFLOW = ROOT / ".github/workflows/rustfs-security-test.yml"
CASE_ROW = "| IAM-101 | user CRUD lifecycle | PASS |"
def named_steps(job: list[str]) -> dict[str, list[str]]:
starts = [i for i, line in enumerate(job) if line.startswith(" - name: ")]
return {
job[start].split(": ", 1)[1].strip('"'): job[start:end]
for start, end in zip(starts, starts[1:] + [len(job)])
}
class SecurityWorkflowTests(unittest.TestCase):
def setUp(self) -> None:
self.source = WORKFLOW.read_text()
self.job = yaml_block(self.source.splitlines(), "security-test", 2)
self.assertIsNotNone(self.job)
starts = [i for i, line in enumerate(self.job) if line.startswith(" - name: ")]
self.steps = {
self.job[start].split(": ", 1)[1].strip('"'): self.job[start:end]
for start, end in zip(starts, starts[1:] + [len(self.job)])
}
self.steps = named_steps(self.job)
self.temp = tempfile.TemporaryDirectory()
self.addCleanup(self.temp.cleanup)
self.directory = Path(self.temp.name)
@@ -192,6 +196,110 @@ class SecurityWorkflowTests(unittest.TestCase):
self.assertNotIn("OLD RUN REPORT", body.read_text())
self.assertIn("https://github.com/rustfs/rustfs/actions/runs/314159", body.read_text())
def test_all_ten_suites_hold_the_shared_lock_for_manual_and_chain_runs(self) -> None:
for suite in ("upgrade", "s3-compat", "kms", "tier", "storage", "heal", "pool-expand", "security", "replication", "performance"):
with self.subTest(suite=suite):
source = (ROOT / f".github/workflows/rustfs-{suite}-test.yml").read_text().splitlines()
# Workflow-level concurrency covers every job, including cleanup,
# regardless of trigger or the runner hosting the job.
self.assertEqual([
line.strip() for line in yaml_block(source, "concurrency", 0)
if line.strip() and not line.lstrip().startswith("#")
], [
"group: rustfs-shared-functional-tests", "cancel-in-progress: false",
])
self.assertIsNotNone(yaml_block(source, "workflow_dispatch", 2))
self.assertIsNotNone(yaml_block(source, "repository_dispatch", 2))
cleanup_name = "Reset test environment (after)" if suite == "performance" else "Cleanup environment (after)"
cleanup = named_steps(yaml_block(source, "jobs", 0))[cleanup_name]
self.assertTrue(any(line.startswith(" if:") and "always()" in line for line in cleanup))
def test_root_dispatches_only_upgrade_and_replication_hands_off_after_failure(self) -> None:
for failed_attempts, issue_exit, token in ((0, 0, "fixture"), (2, 0, "fixture"), (3, 0, "fixture"), (3, 7, "fixture"), (0, 0, "")):
with self.subTest(failed_attempts=failed_attempts, issue_exit=issue_exit, token=bool(token)):
self.setUp()
fake_bin = self.directory / "bin"
fake_bin.mkdir()
commands = {
"gh": '''#!/usr/bin/env bash
set -euo pipefail
if [ "$1" = api ]; then
printf '%s\\n' "$*" >> "$DISPATCHES"
attempt=$(wc -l < "$DISPATCHES")
[ "$attempt" -gt "$FAILED_ATTEMPTS" ]
elif [ "$1 $2" = 'issue create' ]; then
printf 'issue\\n' >> "$EXECUTED"
while [ "$#" -gt 0 ]; do
if [ "$1" = --body-file ]; then
cat "$2" > "$CAPTURE_BODY"
printf '%s\\n' "$2" > "$CAPTURE_BODY_PATH"
fi
shift
done
exit "$ISSUE_EXIT"
else
exit 99
fi
''',
"sleep": '#!/bin/sh\nprintf "sleep %s\\n" "$1" >> "$EXECUTED"\n',
"ssh": '#!/bin/sh\nprintf "cleanup\\n" >> "$EXECUTED"\n',
}
for name, contents in commands.items():
command = fake_bin / name
command.write_text(contents)
command.chmod(0o755)
dispatches = self.directory / "dispatches"
executed = self.directory / "executed"
body = self.directory / "issue-body.md"
body_path = self.directory / "issue-body-path"
self.env.update(
PATH=f"{fake_bin}{os.pathsep}{os.environ['PATH']}", DISPATCHES=str(dispatches),
EXECUTED=str(executed), CAPTURE_BODY=str(body), CAPTURE_BODY_PATH=str(body_path),
FAILED_ATTEMPTS="0", ISSUE_EXIT=str(issue_exit),
RUSTFS_NODES="fixture-node", RUSTFS_SSH_USER="fixture-user",
RUSTFS_NIGHTLY_PACKAGE_URL="https://example.invalid/package.deb",
)
self.context.update({"secrets.PF_TESTING_GH_TOKEN": "fixture", "inputs.suite": "all"})
driver = (ROOT / ".github/workflows/rustfs-functional-chain.yml").read_text()
self.steps = named_steps(yaml_block(driver.splitlines(), "start-chain", 2))
self.assertEqual(list(self.steps), ["Dispatch first suite (upgrade)"])
started = self.run_step("Dispatch first suite (upgrade)")
self.assertEqual(started.returncode, 0, started.stderr)
self.assertEqual(dispatches.read_text().splitlines(), [
"api --method POST repos/rustfs/rustfs/dispatches -f event_type=rustfs-chain-upgrade -F client_payload[from_suite]=nightly-build",
])
dispatches.unlink()
replication = (ROOT / ".github/workflows/rustfs-replication-test.yml").read_text()
job = yaml_block(replication.splitlines(), "replication-test", 2)
self.assertFalse(any(line.startswith(" continue-on-error:") for line in job))
self.steps = named_steps(job)
handoff = "Continue functional chain (next: Performance)"
self.assertIn(" if: ${{ always() && github.event_name == 'repository_dispatch' }}", self.steps[handoff])
self.assertFalse(any(line.strip().startswith("continue-on-error:") for line in self.steps[handoff]))
self.assertIn(" if: always()", self.steps["Cleanup environment (after)"])
self.assertLess(list(self.steps).index("Cleanup environment (after)"), list(self.steps).index(handoff))
suite = self.directory / "auto-testing/rustfs-replication-test.sh"
suite.write_text('#!/bin/sh\nprintf "suite failed\\n" >> "$EXECUTED"\nexit 17\n')
failed = self.run_step("Run replication suite")
self.assertEqual(failed.returncode, 17, failed.stderr)
cleaned = self.run_step("Cleanup environment (after)")
self.assertEqual(cleaned.returncode, 0, cleaned.stderr)
self.assertEqual(executed.read_text().splitlines(), ["suite failed", "cleanup"])
self.env["FAILED_ATTEMPTS"] = str(failed_attempts)
self.context["secrets.PF_TESTING_GH_TOKEN"] = token
forwarded = self.run_step(handoff)
self.assertEqual(forwarded.returncode == 0, bool(token) and failed_attempts < 3, forwarded.stderr)
calls = dispatches.read_text().splitlines() if dispatches.exists() else []
self.assertEqual(calls, [
"api --method POST repos/rustfs/rustfs/dispatches -f event_type=rustfs-chain-performance -F client_payload[from_suite]=replication",
] * (min(failed_attempts + 1, 3) if token else 0))
if failed_attempts == 3:
self.assertIn("could not hand off from **replication** to **Performance**", body.read_text())
self.assertIn("rustfs-chain-performance", body.read_text())
self.assertEqual(executed.read_text().splitlines().count("issue"), 2 if issue_exit else 1)
self.assertFalse(Path(body_path.read_text().strip()).exists())
if __name__ == "__main__":
unittest.main()