fix(release): backport IAM migration startup fixes and recent main fixes (#7825)

This commit is contained in:
Chris
2026-09-14 11:47:13 +08:00
committed by GitHub
parent 95e90dbbc9
commit b61dc6519f
28 changed files with 1135 additions and 148 deletions
+58 -16
View File
@@ -406,13 +406,22 @@ jobs:
# Force rebuild by touching build.rs
touch rustfs/build.rs
if [[ "${{ matrix.cross }}" == "true" ]]; then
# All cross targets in the matrix are Linux; zigbuild handles them.
cargo zigbuild --release --target ${{ matrix.target }} -p rustfs --bin rustfs
else
cargo build --release --target ${{ matrix.target }} -p rustfs --bin rustfs
BINARY_NAMES=(rustfs)
if [[ "${{ matrix.platform }}" == "linux" ]]; then
BINARY_NAMES+=(rustfs-cli)
fi
# Build one binary at a time so release LTO links cannot exhaust a
# self-hosted runner by running concurrently.
for binary in "${BINARY_NAMES[@]}"; do
if [[ "${{ matrix.cross }}" == "true" ]]; then
# All cross targets in the matrix are Linux; zigbuild handles them.
cargo zigbuild --release --target ${{ matrix.target }} -p rustfs --bin "$binary"
else
cargo build --release --target ${{ matrix.target }} -p rustfs --bin "$binary"
fi
done
- name: Create release package
id: package
shell: bash
@@ -498,17 +507,24 @@ jobs:
BINARY_NAME="rustfs"
fi
# Verify the binary exists before packaging
if [[ ! -f "$BINARY_NAME" ]]; then
echo "❌ Binary $BINARY_NAME not found in $(pwd)"
if [[ "${{ matrix.platform }}" == "windows" ]]; then
dir
else
ls -la
fi
exit 1
BINARY_FILES=("$BINARY_NAME")
if [[ "${{ matrix.platform }}" == "linux" ]]; then
BINARY_FILES+=("rustfs-cli")
fi
# Verify the binaries exist before packaging
for binary in "${BINARY_FILES[@]}"; do
if [[ ! -f "$binary" ]]; then
echo "❌ Binary $binary not found in $(pwd)"
if [[ "${{ matrix.platform }}" == "windows" ]]; then
dir
else
ls -la
fi
exit 1
fi
done
# Universal packaging function
package_zip() {
local src=$1
@@ -526,8 +542,12 @@ jobs:
}
# Create the zip package
echo "Start packaging: $BINARY_NAME -> ../../../${PACKAGE_NAME}.zip"
package_zip "$BINARY_NAME" "../../../${PACKAGE_NAME}.zip"
echo "Start packaging: ${BINARY_FILES[*]} -> ../../../${PACKAGE_NAME}.zip"
if [[ "${{ matrix.platform }}" == "linux" ]]; then
zip "../../../${PACKAGE_NAME}.zip" "${BINARY_FILES[@]}"
else
package_zip "$BINARY_NAME" "../../../${PACKAGE_NAME}.zip"
fi
cd ../../..
@@ -601,6 +621,28 @@ jobs:
echo "🔧 Build type: ${BUILD_TYPE}"
echo "📊 Version: ${VERSION}"
- name: Verify packaged Linux CLI
if: matrix.platform == 'linux'
shell: bash
run: |
set -euo pipefail
package_file="${{ steps.package.outputs.package_file }}"
archive_members="$(unzip -Z1 "$package_file" | sort)"
expected_members=$'rustfs\nrustfs-cli'
if [[ "$archive_members" != "$expected_members" ]]; then
echo "Linux package must contain exactly rustfs and rustfs-cli" >&2
printf '%s\n' "$archive_members" >&2
exit 1
fi
if [[ "${{ matrix.target }}" == "x86_64-unknown-linux-gnu" ]]; then
cli_dir="$(mktemp -d)"
trap 'rm -rf "$cli_dir"' EXIT
unzip -q "$package_file" rustfs-cli -d "$cli_dir"
"$cli_dir/rustfs-cli" --help >/dev/null
fi
- name: Verify packaged console
if: matrix.target == 'x86_64-unknown-linux-gnu'
shell: bash
+10
View File
@@ -130,6 +130,16 @@ pub const DEFAULT_SHARD_INTEGRITY_FLEET_CONFIRMED: bool = false;
const _: () = assert!(!DEFAULT_SHARD_INTEGRITY_WRITE);
const _: () = assert!(!DEFAULT_SHARD_INTEGRITY_FLEET_CONFIRMED);
/// How object writes treat a bucket whose stored versioning configuration
/// cannot be parsed: `permissive` writes as if unversioned (the historical
/// behavior, recorded by metrics and an error log) and `strict` refuses the
/// write with 503. Paths that already refuse an unreadable configuration do
/// so in both modes. Any other value fails startup.
pub const ENV_BUCKET_CONFIG_PARSE_MODE: &str = "RUSTFS_BUCKET_CONFIG_PARSE_MODE";
/// Default bucket config parse mode.
pub const DEFAULT_BUCKET_CONFIG_PARSE_MODE: &str = "permissive";
/// Request writing the complete remote-tier version state into object metadata.
///
/// This remains ineffective until
+9 -2
View File
@@ -168,12 +168,19 @@ pub mod bucket {
BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, BUCKET_QUOTA_CONFIG_FILE,
BUCKET_REPLICATION_CONFIG, BUCKET_REQUEST_PAYMENT_CONFIG, BUCKET_SSECONFIG, BUCKET_TABLE_CATALOG_META_PREFIX,
BUCKET_TABLE_CATALOG_TABLE_BUCKETS_PREFIX, BUCKET_TABLE_CONFIG, BUCKET_TABLE_RESERVED_PREFIX, BUCKET_TAGGING_CONFIG,
BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, BUCKET_WEBSITE_CONFIG, BucketMetadata, OBJECT_LOCK_CONFIG,
load_bucket_metadata, table_catalog_path_hash,
BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, BUCKET_WEBSITE_CONFIG, BucketMetadata, ConfigState,
OBJECT_LOCK_CONFIG, UnreadableBucketConfig, is_unreadable_config_error, load_bucket_metadata,
table_catalog_path_hash, unreadable_config_refusal,
};
pub use crate::bucket::metadata::{BUCKET_DURABILITY_CONFIG, BUCKET_ON_DEMAND_MIGRATION_CONFIG};
}
pub mod config_parse_mode {
pub use crate::bucket::config_parse_mode::{
BucketConfigParseMode, bucket_config_parse_mode, validate_bucket_config_parse_mode_env,
};
}
pub mod durability {
pub use crate::bucket::durability::{
BUCKET_DURABILITY_MODE_NONE, BUCKET_DURABILITY_MODE_RELAXED, BUCKET_DURABILITY_MODE_STRICT, BucketDurabilityConfig,
@@ -0,0 +1,241 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Handling of stored bucket sub-configurations whose bytes cannot be parsed
//! (rustfs/backlog#1734): the rollout mode for paths that historically read
//! them as absent, and the metrics that make them visible.
//!
//! The mode only governs those historical degrade paths. Paths that already
//! refuse an unreadable config (the typed getters, the read-modify-write
//! guard, Object Lock and default-encryption decisions, delete-time
//! versioning) refuse in every mode.
use rustfs_config::{DEFAULT_BUCKET_CONFIG_PARSE_MODE, ENV_BUCKET_CONFIG_PARSE_MODE};
use std::collections::HashMap;
use std::sync::{LazyLock, Mutex, OnceLock};
use tracing::error;
/// Counter of stored XML sub-configurations that failed to parse, labeled
/// `bucket`, `config` and `mode`. Increments on every metadata parse, so it
/// must read zero fleet-wide before the default mode flips to strict.
pub const METRIC_BUCKET_METADATA_PARSE_FAILED_TOTAL: &str = "rustfs_bucket_metadata_parse_failed_total";
/// Gauge of buckets whose most recently parsed metadata holds an unreadable
/// XML sub-configuration, labeled `config`. The counter only moves when a
/// parse runs; this reflects the current state between parses.
pub const METRIC_BUCKET_METADATA_UNPARSABLE_CURRENT: &str = "rustfs_bucket_metadata_unparsable_current";
const LOG_COMPONENT: &str = "ecstore";
const LOG_SUBSYSTEM: &str = "bucket_metadata";
const EVENT_CONFIG_UNREADABLE: &str = "bucket_metadata_config_unreadable";
const EVENT_PARSE_MODE_INVALID: &str = "bucket_config_parse_mode_invalid";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BucketConfigParseMode {
/// Historical degrade paths keep reading an unreadable config as absent;
/// the metrics and an error log record every occurrence.
Permissive,
/// Historical degrade paths refuse instead.
Strict,
}
impl BucketConfigParseMode {
pub fn as_str(self) -> &'static str {
match self {
Self::Permissive => "permissive",
Self::Strict => "strict",
}
}
/// Parse a configured value; unset or blank selects the default. An
/// unknown value is an error rather than a silent fallback.
pub fn parse(value: Option<&str>) -> Result<Self, String> {
let value = value
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or(DEFAULT_BUCKET_CONFIG_PARSE_MODE);
match value.to_ascii_lowercase().as_str() {
"permissive" => Ok(Self::Permissive),
"strict" => Ok(Self::Strict),
_ => Err(format!(
"invalid {ENV_BUCKET_CONFIG_PARSE_MODE} value {value:?}; expected permissive or strict"
)),
}
}
}
/// Validate the configured mode; startup calls this so an invalid value fails
/// the node instead of being guessed at.
pub fn validate_bucket_config_parse_mode_env() -> Result<BucketConfigParseMode, String> {
BucketConfigParseMode::parse(rustfs_utils::get_env_opt_str(ENV_BUCKET_CONFIG_PARSE_MODE).as_deref())
}
/// The process-wide mode. Startup has already rejected an invalid value; if
/// this is reached without that validation, an invalid value selects strict,
/// never permissive.
pub fn bucket_config_parse_mode() -> BucketConfigParseMode {
static MODE: OnceLock<BucketConfigParseMode> = OnceLock::new();
*MODE.get_or_init(|| {
validate_bucket_config_parse_mode_env().unwrap_or_else(|err| {
error!(
event = EVENT_PARSE_MODE_INVALID,
component = LOG_COMPONENT,
subsystem = LOG_SUBSYSTEM,
error = %err,
"Invalid bucket config parse mode; using strict"
);
BucketConfigParseMode::Strict
})
})
}
/// Which buckets currently hold which unreadable XML configs, so the gauge
/// reports buckets rather than parse events.
#[derive(Debug, Default)]
struct UnreadableConfigTracker {
by_bucket: HashMap<String, Vec<&'static str>>,
}
impl UnreadableConfigTracker {
/// Replace `bucket`'s unreadable set and return the current bucket count
/// of every config whose count may have changed.
fn update(&mut self, bucket: &str, configs: Vec<&'static str>) -> Vec<(&'static str, usize)> {
let previous = if configs.is_empty() {
self.by_bucket.remove(bucket).unwrap_or_default()
} else {
self.by_bucket.insert(bucket.to_string(), configs.clone()).unwrap_or_default()
};
let mut touched = previous;
touched.extend(configs);
touched.sort_unstable();
touched.dedup();
touched
.into_iter()
.map(|config| (config, self.by_bucket.values().filter(|set| set.contains(&config)).count()))
.collect()
}
}
static TRACKER: LazyLock<Mutex<UnreadableConfigTracker>> = LazyLock::new(Default::default);
fn publish_gauges(bucket: &str, configs: Vec<&'static str>) {
let changed = TRACKER
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.update(bucket, configs);
for (config, count) in changed {
metrics::gauge!(METRIC_BUCKET_METADATA_UNPARSABLE_CURRENT, "config" => config).set(count as f64);
}
}
/// Record the outcome of one metadata parse: every unreadable XML config
/// (config file, stored byte length) counts once, logs at error level, and
/// replaces the bucket's entry in the current-state gauge.
pub(crate) fn record_bucket_config_parse_state(bucket: &str, unreadable: &[(&'static str, usize)]) {
if !unreadable.is_empty() {
let mode = bucket_config_parse_mode();
for (config, _) in unreadable {
metrics::counter!(
METRIC_BUCKET_METADATA_PARSE_FAILED_TOTAL,
"bucket" => bucket.to_string(),
"config" => *config,
"mode" => mode.as_str()
)
.increment(1);
}
error!(
event = EVENT_CONFIG_UNREADABLE,
component = LOG_COMPONENT,
subsystem = LOG_SUBSYSTEM,
bucket = %bucket,
configs = ?unreadable,
mode = mode.as_str(),
"Stored bucket configuration cannot be parsed"
);
}
publish_gauges(bucket, unreadable.iter().map(|(config, _)| *config).collect());
}
/// Drop a deleted bucket from the current-state gauge.
pub(crate) fn forget_bucket_config_parse_state(bucket: &str) {
publish_gauges(bucket, Vec::new());
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bucket::metadata::{BUCKET_TAGGING_CONFIG, BUCKET_VERSIONING_CONFIG};
#[test]
fn parse_mode_defaults_to_permissive_and_rejects_unknown_values() {
assert_eq!(BucketConfigParseMode::parse(None), Ok(BucketConfigParseMode::Permissive));
assert_eq!(BucketConfigParseMode::parse(Some(" ")), Ok(BucketConfigParseMode::Permissive));
assert_eq!(BucketConfigParseMode::parse(Some("permissive")), Ok(BucketConfigParseMode::Permissive));
assert_eq!(BucketConfigParseMode::parse(Some(" Strict ")), Ok(BucketConfigParseMode::Strict));
let err = BucketConfigParseMode::parse(Some("lenient")).expect_err("an unknown mode must not be guessed");
assert!(err.contains(ENV_BUCKET_CONFIG_PARSE_MODE), "{err}");
assert!(err.contains("permissive") && err.contains("strict"), "must list the valid values: {err}");
}
#[test]
fn tracker_counts_buckets_not_parse_events() {
let mut tracker = UnreadableConfigTracker::default();
assert_eq!(tracker.update("a", vec![BUCKET_VERSIONING_CONFIG]), vec![(BUCKET_VERSIONING_CONFIG, 1)]);
// Re-parsing the same state does not double count.
assert_eq!(tracker.update("a", vec![BUCKET_VERSIONING_CONFIG]), vec![(BUCKET_VERSIONING_CONFIG, 1)]);
assert_eq!(tracker.update("b", vec![BUCKET_VERSIONING_CONFIG]), vec![(BUCKET_VERSIONING_CONFIG, 2)]);
// A repaired bucket releases its configs; the changed one is reported.
let mut changed = tracker.update("a", vec![BUCKET_TAGGING_CONFIG]);
changed.sort_unstable();
let mut expected = vec![(BUCKET_TAGGING_CONFIG, 1), (BUCKET_VERSIONING_CONFIG, 1)];
expected.sort_unstable();
assert_eq!(changed, expected);
assert_eq!(tracker.update("b", Vec::new()), vec![(BUCKET_VERSIONING_CONFIG, 0)]);
assert_eq!(tracker.update("never-unreadable", Vec::new()), Vec::new());
}
#[test]
fn every_unreadable_config_increments_the_labeled_counter() {
let recorder = metrics_util::debugging::DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
record_bucket_config_parse_state("metric-bucket", &[(BUCKET_VERSIONING_CONFIG, 12), (BUCKET_TAGGING_CONFIG, 3)]);
record_bucket_config_parse_state("metric-bucket", &[(BUCKET_VERSIONING_CONFIG, 12)]);
record_bucket_config_parse_state("metric-clean-bucket", &[]);
});
let mut versioning = 0;
let mut tagging = 0;
for (composite, _, _, value) in snapshotter.snapshot().into_vec() {
if composite.key().name() != METRIC_BUCKET_METADATA_PARSE_FAILED_TOTAL {
continue;
}
let labels: HashMap<_, _> = composite.key().labels().map(|l| (l.key(), l.value())).collect();
assert_eq!(labels.get("bucket"), Some(&"metric-bucket"));
assert_eq!(labels.get("mode"), Some(&bucket_config_parse_mode().as_str()));
let metrics_util::debugging::DebugValue::Counter(count) = value else {
panic!("parse failures must be a counter");
};
match labels.get("config") {
Some(&config) if config == BUCKET_VERSIONING_CONFIG => versioning += count,
Some(&config) if config == BUCKET_TAGGING_CONFIG => tagging += count,
other => panic!("unexpected config label {other:?}"),
}
}
assert_eq!((versioning, tagging), (2, 1));
}
}
@@ -12691,9 +12691,8 @@ mod tests {
.await
.expect_err("malformed Object Lock metadata must reject lifecycle config resolution");
assert!(
exact_error
.to_string()
.contains("persisted bucket Object Lock configuration is invalid")
crate::bucket::metadata::is_unreadable_config_error(&exact_error),
"malformed Object Lock metadata must surface as the typed unreadable-config refusal: {exact_error}"
);
let runtime_state = install_unconsumed_runtime_expiry_worker(&ecstore, 1).await;
+176 -6
View File
@@ -275,6 +275,141 @@ pub const BUCKET_TABLE_RESERVED_PREFIX: &str = ".rustfs-table";
pub const BUCKET_TABLE_CATALOG_META_PREFIX: &str = "s3tables/catalog";
pub const BUCKET_TABLE_CATALOG_TABLE_BUCKETS_PREFIX: &str = "table-buckets";
/// Refusal to act on a stored sub-configuration whose bytes exist but cannot
/// be parsed. Carried inside [`Error::other`] so callers can tell it apart
/// from a storage fault with [`is_unreadable_config_error`].
#[derive(Debug, Clone)]
pub struct UnreadableBucketConfig {
pub bucket: String,
pub config_file: String,
/// Length of the stored bytes that failed to parse.
pub raw_len: usize,
}
impl std::fmt::Display for UnreadableBucketConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"persisted bucket configuration {} ({} bytes) for bucket {} cannot be parsed; back up the stored bytes (rustfs inspect bucket-meta), then replace or delete the configuration",
self.config_file, self.raw_len, self.bucket
)
}
}
impl std::error::Error for UnreadableBucketConfig {}
pub fn is_unreadable_config_error(err: &Error) -> bool {
unreadable_config_refusal(err).is_some()
}
pub fn unreadable_config_refusal(err: &Error) -> Option<&UnreadableBucketConfig> {
match err {
Error::Io(io) => io.get_ref().and_then(|inner| inner.downcast_ref::<UnreadableBucketConfig>()),
_ => None,
}
}
pub(crate) fn unreadable_config_error(bucket: &str, config_file: &str, raw_len: usize) -> Error {
Error::other(UnreadableBucketConfig {
bucket: bucket.to_string(),
config_file: config_file.to_string(),
raw_len,
})
}
/// Stored state of one sub-configuration as left by
/// [`BucketMetadata::parse_all_configs`]: a parse failure keeps the raw bytes
/// and leaves the typed field `None`, which is kept distinct from "no bytes".
///
/// Deliberately has no `Option` conversion or defaulting accessor: folding
/// `Unreadable` into `Absent` is exactly the failure this type exists to stop.
#[derive(Debug)]
pub enum ConfigState<'a, T> {
/// No bytes are stored.
Absent,
/// The stored bytes parsed.
Valid(&'a T),
/// Bytes are stored but could not be parsed.
Unreadable { raw_len: usize },
}
impl<'a, T> ConfigState<'a, T> {
pub fn of(raw: &[u8], parsed: &'a Option<T>) -> Self {
match parsed {
Some(config) => Self::Valid(config),
None if raw.is_empty() => Self::Absent,
None => Self::Unreadable { raw_len: raw.len() },
}
}
/// `Ok(None)` when absent; a typed [`UnreadableBucketConfig`] refusal
/// when the stored bytes cannot be parsed.
pub fn require(self, bucket: &str, config_file: &str) -> Result<Option<&'a T>> {
match self {
Self::Absent => Ok(None),
Self::Valid(config) => Ok(Some(config)),
Self::Unreadable { raw_len } => Err(unreadable_config_error(bucket, config_file, raw_len)),
}
}
}
/// The XML sub-configurations whose parse failure is retained as
/// [`ConfigState::Unreadable`].
pub const XML_BUCKET_CONFIG_FILES: [&str; 13] = [
BUCKET_NOTIFICATION_CONFIG,
BUCKET_LIFECYCLE_CONFIG,
OBJECT_LOCK_CONFIG,
BUCKET_VERSIONING_CONFIG,
BUCKET_SSECONFIG,
BUCKET_TAGGING_CONFIG,
BUCKET_REPLICATION_CONFIG,
BUCKET_CORS_CONFIG,
BUCKET_LOGGING_CONFIG,
BUCKET_WEBSITE_CONFIG,
BUCKET_ACCELERATE_CONFIG,
BUCKET_REQUEST_PAYMENT_CONFIG,
BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG,
];
impl BucketMetadata {
/// Whether the stored XML sub-configuration `config_file` has bytes that
/// cannot be parsed. Non-XML config files report `false`: they carry their
/// own unreadable handling (policy, quota, bucket targets).
pub fn xml_config_unreadable(&self, config_file: &str) -> bool {
self.xml_config_unreadable_len(config_file).is_some()
}
/// Length of the stored bytes of XML sub-configuration `config_file` when
/// they cannot be parsed; `None` when it is absent, readable, or not an
/// XML config.
pub fn xml_config_unreadable_len(&self, config_file: &str) -> Option<usize> {
fn unreadable<T>(raw: &[u8], parsed: &Option<T>) -> Option<usize> {
match ConfigState::of(raw, parsed) {
ConfigState::Unreadable { raw_len } => Some(raw_len),
ConfigState::Absent | ConfigState::Valid(_) => None,
}
}
match config_file {
BUCKET_NOTIFICATION_CONFIG => unreadable(&self.notification_config_xml, &self.notification_config),
BUCKET_LIFECYCLE_CONFIG => unreadable(&self.lifecycle_config_xml, &self.lifecycle_config),
OBJECT_LOCK_CONFIG => unreadable(&self.object_lock_config_xml, &self.object_lock_config),
BUCKET_VERSIONING_CONFIG => unreadable(&self.versioning_config_xml, &self.versioning_config),
BUCKET_SSECONFIG => unreadable(&self.encryption_config_xml, &self.sse_config),
BUCKET_TAGGING_CONFIG => unreadable(&self.tagging_config_xml, &self.tagging_config),
BUCKET_REPLICATION_CONFIG => unreadable(&self.replication_config_xml, &self.replication_config),
BUCKET_CORS_CONFIG => unreadable(&self.cors_config_xml, &self.cors_config),
BUCKET_LOGGING_CONFIG => unreadable(&self.logging_config_xml, &self.logging_config),
BUCKET_WEBSITE_CONFIG => unreadable(&self.website_config_xml, &self.website_config),
BUCKET_ACCELERATE_CONFIG => unreadable(&self.accelerate_config_xml, &self.accelerate_config),
BUCKET_REQUEST_PAYMENT_CONFIG => unreadable(&self.request_payment_config_xml, &self.request_payment_config),
BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG => {
unreadable(&self.public_access_block_config_xml, &self.public_access_block_config)
}
_ => None,
}
}
}
pub fn table_catalog_path_hash(value: &str) -> String {
let digest = Sha256::digest(value.as_bytes());
let mut output = String::with_capacity(digest.len() * 2);
@@ -473,6 +608,13 @@ impl BucketMetadata {
self.lock_enabled || self.object_lock_config.as_ref().is_some_and(|v| v.enabled())
}
/// Whether an operation that may skip Object Lock checks must keep them.
/// Stored lock bytes that cannot be parsed leave the lock state unknown,
/// which must not be read as "no Object Lock".
pub fn object_lock_checks_required(&self) -> bool {
self.object_locking() || self.xml_config_unreadable(OBJECT_LOCK_CONFIG)
}
pub fn table_bucket_enabled(&self) -> bool {
!self.table_bucket_config_json.is_empty()
}
@@ -985,17 +1127,25 @@ impl BucketMetadata {
/// |---|---|
/// | policy | Fails closed: `get_bucket_policy` re-parses the raw JSON and propagates the error; `get_bucket_policy_raw` returns the stored bytes. |
/// | object lock | Fails closed in `object_lock_config_state_from_authoritative_metadata`; a retention decision may never be taken on a guess. |
/// | versioning | Fails closed in `get_versioning_config`; guessing Unversioned would make delete markers and version ids diverge from what is on disk. |
/// | versioning | Fails closed in `get_versioning_config` and the delete-time snapshot; guessing Unversioned would make delete markers and version ids diverge from what is on disk. Object writes lay out versions through `BucketVersioningSys::get_for_write`, which refuses in strict mode and, in the default permissive mode, keeps the historical unversioned write (see `config_parse_mode`). |
/// | replication | Fails closed in `get_replication_config`. |
/// | bucket targets | Fails closed in `get_bucket_targets_config`, and `sync_bucket_target_sys` marks the bucket unreadable in `BucketTargetSys` instead of publishing an empty target set (rustfs/backlog#2282). |
/// | encryption | Fails closed in `get_sse_config`: degrading to "no default encryption" stores plaintext objects the operator required to be encrypted. |
/// | public access block | Fails closed in `get_public_access_block_config`: degrading grants the anonymous access the operator asked to block. |
/// | quota | Fails closed in `get_quota_config`; the enforcement path in `quota::checker` already re-parses the raw JSON and refuses on error. |
/// | lifecycle | Safe to degrade: no rules means no expiration and no transition, so nothing is deleted or moved on the strength of an unreadable rule set. The bucket keeps serving reads and writes. |
/// | notification | Safe to degrade: events are an outbound side channel; no consumer draws a durability or authorization conclusion from their absence. |
/// | tagging | Safe to degrade: bucket tags are cost-allocation labels here; object-level tag conditions come from object metadata, not this blob. |
/// | CORS | Safe to degrade: an absent CORS configuration rejects cross-origin browser requests, which is already the restrictive direction. |
/// | logging, website, accelerate, request payment, bucket ACL | Safe to degrade: each only shapes an optional response or an optional side channel, and none of them authorizes an action or decides whether data is retained. |
/// | lifecycle | `get_lifecycle_config` fails closed, so GetBucketLifecycle reports the fault instead of NoSuchLifecycleConfiguration. ILM, scanner and expiry-header consumers still degrade to "no rules": nothing is deleted or moved on the strength of an unreadable rule set, and the bucket keeps serving reads and writes. |
/// | notification | `get_notification_config` fails closed, so GetBucketNotificationConfiguration reports the fault and startup leaves that one bucket's rules unchanged instead of clearing them; other buckets are unaffected. |
/// | tagging, CORS, logging, website, accelerate, request payment | The getter fails closed so the matching GET API reports the fault instead of "not configured". Per-request consumers (CORS response headers) still degrade to "not configured", which is the restrictive direction. |
/// | bucket ACL | Safe to degrade: it only shapes an optional response. |
///
/// Independently of the table, `update_config_with` refuses a
/// read-modify-write of an unreadable XML config before its mutate step
/// runs, so the unreadable bytes are never replaced by a rewrite that saw
/// them as absent. Other configs of the same bucket stay writable.
///
/// Every unreadable XML config found here is counted in
/// `rustfs_bucket_metadata_parse_failed_total`, reflected in
/// `rustfs_bucket_metadata_unparsable_current`, and logged at error level.
pub(super) fn parse_all_configs(&mut self) -> Result<()> {
if let Err(e) = self.parse_policy_config() {
tracing::warn!(
@@ -1239,6 +1389,14 @@ impl BucketMetadata {
);
}
if !self.name.is_empty() {
let unreadable: Vec<(&'static str, usize)> = XML_BUCKET_CONFIG_FILES
.iter()
.filter_map(|config| self.xml_config_unreadable_len(config).map(|raw_len| (*config, raw_len)))
.collect();
super::config_parse_mode::record_bucket_config_parse_state(&self.name, &unreadable);
}
Ok(())
}
}
@@ -1356,6 +1514,18 @@ where
mod test {
use super::*;
/// rustfs/backlog#1734: `StorageError::clone` rebuilds I/O errors from
/// their text. The unreadable-config refusal must stay typed across that
/// clone, or a cloned error stops mapping to its retryable S3 response.
#[test]
fn unreadable_config_refusal_survives_storage_error_clone() {
let err = unreadable_config_error("b", BUCKET_TAGGING_CONFIG, 7);
let cloned = err.clone();
assert!(is_unreadable_config_error(&cloned), "clone lost the typed refusal: {cloned}");
assert_eq!(unreadable_config_refusal(&cloned).map(|r| r.raw_len), Some(7));
assert!(is_unreadable_config_error(&err));
}
/// Decode a whitespace-tolerant hex fixture into bytes.
fn decode_hex(s: &str) -> Vec<u8> {
let s: String = s.chars().filter(|c| !c.is_whitespace()).collect();
+255 -64
View File
@@ -13,7 +13,9 @@
// limitations under the License.
use super::metadata::{
BUCKET_TARGETS_FILE, BucketMetadata, load_bucket_incarnation, load_bucket_metadata, save_bucket_incarnation,
BUCKET_ACCELERATE_CONFIG, BUCKET_CORS_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_LOGGING_CONFIG, BUCKET_REQUEST_PAYMENT_CONFIG,
BUCKET_TAGGING_CONFIG, BUCKET_TARGETS_FILE, BUCKET_WEBSITE_CONFIG, BucketMetadata, ConfigState, load_bucket_incarnation,
load_bucket_metadata, save_bucket_incarnation, unreadable_config_error,
};
use super::quota::BucketQuota;
use super::target::BucketTargets;
@@ -150,8 +152,8 @@ enum BucketMetadataAuthority {
}
pub(crate) fn object_lock_config_state_from_authoritative_metadata(bm: &BucketMetadata) -> Result<ObjectLockConfigState> {
if bm.object_lock_config.is_none() && !bm.object_lock_config_xml.is_empty() {
return Err(Error::other("persisted bucket Object Lock configuration is invalid"));
if let Some(raw_len) = bm.xml_config_unreadable_len(super::metadata::OBJECT_LOCK_CONFIG) {
return Err(unreadable_config_error(&bm.name, super::metadata::OBJECT_LOCK_CONFIG, raw_len));
}
if let Some(config) = bm.object_lock_config.clone() {
@@ -1863,6 +1865,7 @@ impl BucketMetadataSys {
drop(map);
let removed_fabricated = self.fabricated_metadata.write().await.remove(bucket);
self.missing_buckets.insert(bucket.to_string(), ()).await;
super::config_parse_mode::forget_bucket_config_parse_state(bucket);
if removed {
BucketTargetSys::get().delete(bucket).await;
clear_bucket_durability(bucket);
@@ -1956,6 +1959,13 @@ impl BucketMetadataSys {
if !bm.bucket_incarnation_sidecar || bm.bucket_incarnation_id != expected_incarnation_id {
return Err(Error::BucketNotFound(bucket.to_string()));
}
// `mutate` would see an unreadable config as absent and rebuild it
// from nothing; persisting that destroys the only copy of the stored
// bytes. Only the rewritten config is checked: `update_config` carries
// every other raw config through unchanged.
if let Some(raw_len) = bm.xml_config_unreadable_len(config_file) {
return Err(unreadable_config_error(bucket, config_file, raw_len));
}
let data = mutate(&bm)?;
let updated = bm.update_config(config_file, data)?;
@@ -2198,12 +2208,11 @@ impl BucketMetadataSys {
}
};
if !bm.versioning_config_xml.is_empty() && bm.versioning_config.is_none() {
Err(Error::other("persisted bucket versioning configuration is invalid"))
} else if let Some(config) = &bm.versioning_config {
Ok((config.clone(), bm.versioning_config_updated_at))
} else {
Ok((VersioningConfiguration::default(), bm.versioning_config_updated_at))
match ConfigState::of(&bm.versioning_config_xml, &bm.versioning_config)
.require(bucket, super::metadata::BUCKET_VERSIONING_CONFIG)?
{
Some(config) => Ok((config.clone(), bm.versioning_config_updated_at)),
None => Ok((VersioningConfiguration::default(), bm.versioning_config_updated_at)),
}
}
@@ -2212,8 +2221,8 @@ impl BucketMetadataSys {
return Ok(false);
};
if metadata.versioning_config.is_none() && !metadata.versioning_config_xml.is_empty() {
return Err(Error::other("persisted bucket versioning configuration is invalid"));
if let Some(raw_len) = metadata.xml_config_unreadable_len(super::metadata::BUCKET_VERSIONING_CONFIG) {
return Err(unreadable_config_error(bucket, super::metadata::BUCKET_VERSIONING_CONFIG, raw_len));
}
Ok(metadata.versioning_config.is_none() && metadata.versioning_config_xml.is_empty())
@@ -2274,22 +2283,20 @@ impl BucketMetadataSys {
pub async fn get_tagging_config(&self, bucket: &str) -> Result<(Tagging, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.tagging_config {
Ok((config.clone(), bm.tagging_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.tagging_config_xml, &bm.tagging_config).require(bucket, BUCKET_TAGGING_CONFIG)? {
Some(config) => Ok((config.clone(), bm.tagging_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
pub async fn get_public_access_block_config(&self, bucket: &str) -> Result<(PublicAccessBlockConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if !bm.public_access_block_config_xml.is_empty() && bm.public_access_block_config.is_none() {
Err(Error::other("persisted bucket public access block configuration is invalid"))
} else if let Some(config) = &bm.public_access_block_config {
Ok((config.clone(), bm.public_access_block_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.public_access_block_config_xml, &bm.public_access_block_config)
.require(bucket, super::metadata::BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG)?
{
Some(config) => Ok((config.clone(), bm.public_access_block_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
@@ -2571,91 +2578,79 @@ impl BucketMetadataSys {
pub async fn get_lifecycle_config(&self, bucket: &str) -> Result<(BucketLifecycleConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.lifecycle_config {
if config.rules.is_empty() {
Err(Error::ConfigNotFound)
} else {
Ok((config.clone(), bm.lifecycle_config_updated_at))
}
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.lifecycle_config_xml, &bm.lifecycle_config).require(bucket, BUCKET_LIFECYCLE_CONFIG)? {
Some(config) if !config.rules.is_empty() => Ok((config.clone(), bm.lifecycle_config_updated_at)),
_ => Err(Error::ConfigNotFound),
}
}
pub async fn get_notification_config(&self, bucket: &str) -> Result<Option<NotificationConfiguration>> {
let bm = match self.get_config(bucket).await {
Ok((bm, _)) => bm.notification_config.clone(),
Err(err) => {
if err == Error::ConfigNotFound {
None
} else {
return Err(err);
}
}
Ok((bm, _)) => bm,
Err(Error::ConfigNotFound) => return Ok(None),
Err(err) => return Err(err),
};
Ok(bm)
// Unreadable must not read as "no notification configured": that
// would silently drop the bucket's event rules.
Ok(ConfigState::of(&bm.notification_config_xml, &bm.notification_config)
.require(bucket, super::metadata::BUCKET_NOTIFICATION_CONFIG)?
.cloned())
}
pub async fn get_sse_config(&self, bucket: &str) -> Result<(ServerSideEncryptionConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if !bm.encryption_config_xml.is_empty() && bm.sse_config.is_none() {
Err(Error::other("persisted bucket encryption configuration is invalid"))
} else if let Some(config) = &bm.sse_config {
Ok((config.clone(), bm.encryption_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.encryption_config_xml, &bm.sse_config).require(bucket, super::metadata::BUCKET_SSECONFIG)? {
Some(config) => Ok((config.clone(), bm.encryption_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
pub async fn get_cors_config(&self, bucket: &str) -> Result<(CORSConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.cors_config {
Ok((config.clone(), bm.cors_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.cors_config_xml, &bm.cors_config).require(bucket, BUCKET_CORS_CONFIG)? {
Some(config) => Ok((config.clone(), bm.cors_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
pub async fn get_website_config(&self, bucket: &str) -> Result<(WebsiteConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.website_config {
Ok((config.clone(), bm.website_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.website_config_xml, &bm.website_config).require(bucket, BUCKET_WEBSITE_CONFIG)? {
Some(config) => Ok((config.clone(), bm.website_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
pub async fn get_logging_config(&self, bucket: &str) -> Result<(BucketLoggingStatus, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.logging_config {
Ok((config.clone(), bm.logging_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.logging_config_xml, &bm.logging_config).require(bucket, BUCKET_LOGGING_CONFIG)? {
Some(config) => Ok((config.clone(), bm.logging_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
pub async fn get_accelerate_config(&self, bucket: &str) -> Result<(AccelerateConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.accelerate_config {
Ok((config.clone(), bm.accelerate_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.accelerate_config_xml, &bm.accelerate_config).require(bucket, BUCKET_ACCELERATE_CONFIG)? {
Some(config) => Ok((config.clone(), bm.accelerate_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
pub async fn get_request_payment_config(&self, bucket: &str) -> Result<(RequestPaymentConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
if let Some(config) = &bm.request_payment_config {
Ok((config.clone(), bm.request_payment_config_updated_at))
} else {
Err(Error::ConfigNotFound)
match ConfigState::of(&bm.request_payment_config_xml, &bm.request_payment_config)
.require(bucket, BUCKET_REQUEST_PAYMENT_CONFIG)?
{
Some(config) => Ok((config.clone(), bm.request_payment_config_updated_at)),
None => Err(Error::ConfigNotFound),
}
}
@@ -3826,6 +3821,202 @@ mod tests {
assert!(matches!(err, Error::Io(_)), "malformed persisted policy must surface its parse failure");
}
/// Persist `bucket` with the given raw sub-configuration bytes, bypassing
/// the parse step the way a newer or foreign writer (or disk damage) would.
async fn persist_bucket_with_raw_config(
sys: &BucketMetadataSys,
dirs: &[tempfile::TempDir],
bucket: &str,
config_file: &str,
raw: &[u8],
) {
for dir in dirs {
std::fs::create_dir_all(dir.path().join(bucket)).expect("bucket volume should be created");
}
let mut bm = BucketMetadata::new(bucket);
bm.update_config(config_file, raw.to_vec())
.expect("raw config should be stored");
sys.persist_new_and_set(bm).await.expect("initial metadata should persist");
}
/// rustfs/backlog#1734: a read-modify-write of a stored config that cannot
/// be parsed must be refused before `mutate` runs. Otherwise `mutate` sees
/// the unreadable config as absent, rebuilds it from nothing, and the
/// write-back destroys the only copy of the original bytes.
#[tokio::test]
async fn update_config_with_refuses_rewrite_of_unreadable_target_config() {
use crate::bucket::metadata::BUCKET_TAGGING_CONFIG;
let (dirs, ecstore) = isolated_store_over_temp_disks().await;
let sys = BucketMetadataSys::new(ecstore);
let bucket = "unreadable-tagging-rmw";
let corrupt = b"<Tagging><TagSet><Tag><Key>team</Key>".to_vec();
persist_bucket_with_raw_config(&sys, &dirs, bucket, BUCKET_TAGGING_CONFIG, &corrupt).await;
let mutate_calls = std::sync::atomic::AtomicUsize::new(0);
let err = sys
.update_config_with(bucket, BUCKET_TAGGING_CONFIG, |_| {
mutate_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Ok(Vec::new())
})
.await
.expect_err("a rewrite of an unreadable config must be refused");
assert_eq!(
mutate_calls.load(std::sync::atomic::Ordering::SeqCst),
0,
"mutate must not see an unreadable config as absent"
);
assert!(
crate::bucket::metadata::is_unreadable_config_error(&err),
"the refusal must be identifiable as an unreadable-config refusal: {err}"
);
sys.metadata_map.write().await.clear();
let (reloaded, _) = sys.get_config(bucket).await.expect("metadata should reload from disk");
assert_eq!(reloaded.tagging_config_xml, corrupt, "the original bytes must stay untouched");
}
/// rustfs/backlog#1734: the refusal is per config. One unreadable config
/// must not block a read-modify-write of a different, readable config, and
/// that write must carry the unreadable bytes through unchanged.
#[tokio::test]
async fn update_config_with_rewrites_readable_config_beside_unreadable_one() {
use crate::bucket::metadata::{BUCKET_CORS_CONFIG, BUCKET_TAGGING_CONFIG};
let (dirs, ecstore) = isolated_store_over_temp_disks().await;
let sys = BucketMetadataSys::new(ecstore);
let bucket = "unreadable-tagging-cors-rmw";
let corrupt = b"<Tagging><TagSet><Tag><Key>team</Key>".to_vec();
persist_bucket_with_raw_config(&sys, &dirs, bucket, BUCKET_TAGGING_CONFIG, &corrupt).await;
let xml = br#"<CORSConfiguration><CORSRule><AllowedMethod>GET</AllowedMethod><AllowedOrigin>https://example.test</AllowedOrigin></CORSRule></CORSConfiguration>"#.to_vec();
sys.update_config_with(bucket, BUCKET_CORS_CONFIG, move |_| Ok(xml))
.await
.expect("a readable config must stay writable beside an unreadable one");
sys.metadata_map.write().await.clear();
let (reloaded, _) = sys.get_config(bucket).await.expect("metadata should reload from disk");
assert_eq!(
reloaded.tagging_config_xml, corrupt,
"the unreadable config must be carried through byte-for-byte"
);
let (stored_cors, _) = sys.get_cors_config(bucket).await.expect("cors should be readable");
assert_eq!(stored_cors.cors_rules.len(), 1);
}
/// rustfs/backlog#1734: reading a stored config that cannot be parsed
/// must fail, not report the config as absent (which the S3 GET handlers
/// turn into NoSuchTagSet / NoSuchLifecycleConfiguration).
#[tokio::test]
async fn unreadable_tagging_and_lifecycle_reads_fail_instead_of_reading_absent() {
use crate::bucket::metadata::{BUCKET_LIFECYCLE_CONFIG, BUCKET_TAGGING_CONFIG};
let (dirs, ecstore) = isolated_store_over_temp_disks().await;
let sys = BucketMetadataSys::new(ecstore);
let tagging_bucket = "unreadable-tagging-read";
persist_bucket_with_raw_config(&sys, &dirs, tagging_bucket, BUCKET_TAGGING_CONFIG, b"<Tagging><TagSet>").await;
let err = sys
.get_tagging_config(tagging_bucket)
.await
.expect_err("unreadable tagging must not read as a value");
assert_ne!(err, Error::ConfigNotFound, "unreadable tagging must not read as absent");
let lifecycle_bucket = "unreadable-lifecycle-read";
persist_bucket_with_raw_config(&sys, &dirs, lifecycle_bucket, BUCKET_LIFECYCLE_CONFIG, b"<LifecycleConfiguration><Rule>")
.await;
let err = sys
.get_lifecycle_config(lifecycle_bucket)
.await
.expect_err("unreadable lifecycle must not read as a value");
assert_ne!(err, Error::ConfigNotFound, "unreadable lifecycle must not read as absent");
// The genuinely absent case still reads as absent.
let absent_bucket = "absent-tagging-read-control";
persist_bucket_with_raw_config(&sys, &dirs, absent_bucket, BUCKET_TAGGING_CONFIG, b"").await;
assert_eq!(sys.get_tagging_config(absent_bucket).await.expect_err("absent"), Error::ConfigNotFound);
}
/// rustfs/backlog#1734: the configs that gate object writes and deletes
/// (versioning, Object Lock, default encryption) and the notification
/// config must refuse with the typed unreadable-config error, so the S3
/// layer can answer 503 with the bucket and config named instead of a
/// generic 500, and notification setup can isolate the one bucket.
#[tokio::test]
async fn unreadable_gating_configs_refuse_with_the_typed_error() {
use crate::bucket::metadata::{BUCKET_VERSIONING_CONFIG, unreadable_config_refusal};
let (dirs, ecstore) = isolated_store_over_temp_disks().await;
let sys = BucketMetadataSys::new(ecstore);
// `update_config` validates some configs on write, so the corrupt
// bytes go straight into the raw fields, as a damaged object would.
let persist_corrupt = |bucket: &'static str, corrupt: fn(&mut BucketMetadata)| {
for dir in &dirs {
std::fs::create_dir_all(dir.path().join(bucket)).expect("bucket volume should be created");
}
let mut bm = BucketMetadata::new(bucket);
corrupt(&mut bm);
sys.persist_new_and_set(bm)
};
persist_corrupt("unreadable-versioning", |bm| {
bm.versioning_config_xml = b"<VersioningConfiguration><Status>".to_vec();
})
.await
.expect("corrupt versioning should persist");
let err = sys
.get_versioning_config("unreadable-versioning")
.await
.expect_err("unreadable versioning must not read as a value");
let refusal = unreadable_config_refusal(&err).unwrap_or_else(|| panic!("expected a typed refusal, got {err}"));
assert_eq!(refusal.config_file, BUCKET_VERSIONING_CONFIG);
assert_eq!(refusal.raw_len, b"<VersioningConfiguration><Status>".len());
persist_corrupt("unreadable-lock", |bm| bm.object_lock_config_xml = b"<ObjectLockConfiguration>".to_vec())
.await
.expect("corrupt Object Lock should persist");
let err = sys
.get_object_lock_config_state("unreadable-lock")
.await
.expect_err("unreadable Object Lock must not read as a value");
assert!(unreadable_config_refusal(&err).is_some(), "{err}");
persist_corrupt("unreadable-sse", |bm| {
bm.encryption_config_xml = b"<ServerSideEncryptionConfiguration>".to_vec();
})
.await
.expect("corrupt encryption should persist");
let err = sys
.get_sse_config("unreadable-sse")
.await
.expect_err("unreadable encryption must not read as a value");
assert!(unreadable_config_refusal(&err).is_some(), "{err}");
persist_corrupt("unreadable-notify", |bm| {
bm.notification_config_xml = b"<NotificationConfiguration>".to_vec();
})
.await
.expect("corrupt notification should persist");
let err = sys
.get_notification_config("unreadable-notify")
.await
.expect_err("unreadable notification must not read as \"no notification configured\"");
assert!(unreadable_config_refusal(&err).is_some(), "{err}");
// Absent stays absent.
persist_corrupt("absent-notification-read", |_| {})
.await
.expect("plain bucket should persist");
assert!(
sys.get_notification_config("absent-notification-read")
.await
.expect("absent")
.is_none()
);
}
/// A tagging rewrite through `update_config_with` (the Swift metadata
/// POST path) is persisted: it survives a metadata reload from disk, and
/// an emptied rewrite clears the config in the cached copy too instead of
+1
View File
@@ -16,6 +16,7 @@
pub mod bandwidth;
pub mod bucket_target_sys;
pub mod config_parse_mode;
pub mod durability;
pub mod error;
pub mod lifecycle;
@@ -179,9 +179,11 @@ fn replication_config_from_metadata(metadata: &BucketMetadata) -> Result<Option<
}
fn delete_request_snapshot_from_metadata(metadata: Arc<BucketMetadata>) -> Result<DeleteReplicationConfigSnapshot> {
if !metadata.versioning_config_xml.is_empty() && metadata.versioning_config.is_none() {
return Err(super::replication_error_boundary::Error::other(
"persisted bucket versioning configuration is invalid",
if let Some(raw_len) = metadata.xml_config_unreadable_len(crate::bucket::metadata::BUCKET_VERSIONING_CONFIG) {
return Err(crate::bucket::metadata::unreadable_config_error(
&metadata.name,
crate::bucket::metadata::BUCKET_VERSIONING_CONFIG,
raw_len,
));
}
+104 -1
View File
@@ -12,11 +12,13 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use super::config_parse_mode::{BucketConfigParseMode, bucket_config_parse_mode};
use super::metadata::unreadable_config_refusal;
use super::{metadata_sys::get_bucket_metadata_sys, versioning::VersioningApi};
use crate::disk::RUSTFS_META_BUCKET;
use crate::error::Result;
use s3s::dto::VersioningConfiguration;
use tracing::warn;
use tracing::{error, warn};
pub struct BucketVersioningSys {}
@@ -86,6 +88,23 @@ impl BucketVersioningSys {
Ok(cfg)
}
/// Versioning configuration for laying out an object write.
///
/// An unreadable stored configuration is refused in strict mode: writing
/// as if unversioned would overwrite versions of a bucket that may have
/// versioning enabled. Any other lookup failure keeps the historical
/// fallback to the default configuration.
pub async fn get_for_write(bucket: &str) -> Result<VersioningConfiguration> {
resolve_versioning_for_write(bucket, Self::get(bucket).await, bucket_config_parse_mode())
}
/// `(versioned, version_suspended)` for an object write under `prefix`,
/// from one [`Self::get_for_write`] lookup.
pub async fn write_state(bucket: &str, prefix: &str) -> Result<(bool, bool)> {
let config = Self::get_for_write(bucket).await?;
Ok((config.prefix_enabled(prefix), config.prefix_suspended(prefix)))
}
/// Instance-scoped variant of [`Self::get`] (backlog#1052): resolves the
/// caller's own instance context so a second in-process store never
/// answers with the first instance's versioning state; falls back to the
@@ -103,3 +122,87 @@ impl BucketVersioningSys {
Ok(cfg)
}
}
fn resolve_versioning_for_write(
bucket: &str,
lookup: Result<VersioningConfiguration>,
mode: BucketConfigParseMode,
) -> Result<VersioningConfiguration> {
match lookup {
Ok(config) => Ok(config),
Err(err) => match (unreadable_config_refusal(&err), mode) {
(Some(_), BucketConfigParseMode::Strict) => Err(err),
(Some(refusal), BucketConfigParseMode::Permissive) => {
// RUSTFS_COMPAT_TODO(s3gate-parse-strict): permissive mode keeps the historical unversioned write for an unreadable versioning config so a rollout can measure affected buckets first. Remove after the parse-failure metric has read zero fleet-wide for two releases and one further release has shipped with strict as the default.
error!(
event = "bucket_versioning_config_unreadable",
component = "ecstore",
subsystem = "bucket_versioning",
bucket = %bucket,
config = %refusal.config_file,
raw_len = refusal.raw_len,
mode = mode.as_str(),
result = "write_unversioned",
"Bucket versioning configuration is unreadable; writing as unversioned"
);
Ok(VersioningConfiguration::default())
}
(None, _) => {
warn!(bucket = %bucket, error = ?err, "failed to load bucket versioning configuration; using default configuration");
Ok(VersioningConfiguration::default())
}
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bucket::metadata::{BUCKET_VERSIONING_CONFIG, is_unreadable_config_error, unreadable_config_error};
use crate::error::Error;
fn enabled() -> VersioningConfiguration {
VersioningConfiguration {
status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::ENABLED)),
..Default::default()
}
}
/// rustfs/backlog#1734: strict mode refuses a write whose versioning
/// state is unknown instead of laying it out as unversioned.
#[test]
fn strict_mode_refuses_a_write_against_unreadable_versioning() {
let err = resolve_versioning_for_write(
"b",
Err(unreadable_config_error("b", BUCKET_VERSIONING_CONFIG, 9)),
BucketConfigParseMode::Strict,
)
.expect_err("strict mode must refuse");
assert!(is_unreadable_config_error(&err), "{err}");
}
#[test]
fn permissive_mode_keeps_the_historical_unversioned_write() {
let config = resolve_versioning_for_write(
"b",
Err(unreadable_config_error("b", BUCKET_VERSIONING_CONFIG, 9)),
BucketConfigParseMode::Permissive,
)
.expect("permissive mode keeps writing");
assert!(!config.enabled());
}
#[test]
fn readable_config_and_lookup_faults_behave_as_before_in_every_mode() {
for mode in [BucketConfigParseMode::Permissive, BucketConfigParseMode::Strict] {
assert!(
resolve_versioning_for_write("b", Ok(enabled()), mode)
.expect("readable")
.enabled()
);
let fallback =
resolve_versioning_for_write("b", Err(Error::other("metadata read failed")), mode).expect("historical fallback");
assert!(!fallback.enabled());
}
}
}
+4
View File
@@ -636,6 +636,10 @@ impl Clone for StorageError {
source: Box::new(context.clone()),
},
))
} else if let Some(refusal) = crate::bucket::metadata::unreadable_config_refusal(self) {
// Keep the refusal typed so a cloned error still maps to
// its retryable S3 response (rustfs/backlog#1734).
StorageError::Io(std::io::Error::new(e.kind(), refusal.clone()))
} else {
StorageError::Io(std::io::Error::new(e.kind(), e.to_string()))
}
+10 -1
View File
@@ -5659,7 +5659,7 @@ fn check_object_lock_retention_update(bucket: &str, object: &str, obj_info: &Obj
/// object-lock protection is never skipped because of a metadata lookup miss.
#[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")]
pub(crate) fn object_lock_delete_check_required(bucket_meta: Option<&crate::bucket::metadata::BucketMetadata>) -> bool {
bucket_meta.is_none_or(|meta| meta.object_locking())
bucket_meta.is_none_or(|meta| meta.object_lock_checks_required())
}
fn restore_expiry_snapshot_matches(obj_info: &ObjectInfo, opts: &ObjectOptions) -> bool {
@@ -12140,6 +12140,15 @@ mod tests {
assert!(object_lock_delete_check_required(None));
}
/// rustfs/backlog#1734: stored Object Lock bytes that cannot be parsed
/// mean the lock state is unknown, not absent; the check must stay on.
#[test]
fn test_object_lock_delete_check_required_fails_closed_on_unreadable_lock_config() {
let mut bm = crate::bucket::metadata::BucketMetadata::new("unreadable-lock-bucket");
bm.object_lock_config_xml = b"<ObjectLockConfiguration>".to_vec();
assert!(object_lock_delete_check_required(Some(&bm)));
}
#[test]
fn test_should_persist_encryption_original_size_rejects_plain_metadata() {
let metadata = HashMap::from([("content-type".to_string(), "application/octet-stream".to_string())]);
+6 -15
View File
@@ -15,8 +15,8 @@
use std::collections::HashMap;
use std::sync::Arc;
use rustfs_ecstore::api::bucket::metadata::BUCKET_TAGGING_CONFIG;
pub(crate) use rustfs_ecstore::api::bucket::metadata::BucketMetadata as SwiftBucketMetadata;
use rustfs_ecstore::api::bucket::metadata::{BUCKET_TAGGING_CONFIG, is_unreadable_config_error};
use rustfs_ecstore::api::bucket::metadata_sys::{get as get_swift_bucket_metadata_from_backend, update_config_with};
use rustfs_ecstore::api::bucket::utils::serialize as serialize_bucket_config;
use rustfs_ecstore::api::error::Error as SwiftStorageError;
@@ -62,11 +62,6 @@ const LOG_COMPONENT_PROTOCOLS: &str = "protocols";
const LOG_SUBSYSTEM_SWIFT_STORAGE: &str = "swift_storage";
const EVENT_SWIFT_BUCKET_TAGGING_UPDATE: &str = "swift_bucket_tagging_update";
/// Marks the refusal to rewrite an unreadable persisted tagging config, so the
/// caller can turn it into an actionable client error rather than a generic
/// storage failure. Carried through the ecstore error, which is a string type.
const UNREADABLE_TAGGING_SENTINEL: &str = "swift: persisted tagging config could not be parsed";
pub type SwiftGetObjectReader = <SwiftStore as storage_contracts::ObjectIO>::GetObjectReader;
pub type SwiftObjectInfo = <SwiftStore as storage_contracts::ObjectOperations>::ObjectInfo;
pub type SwiftObjectOptions = <SwiftStore as storage_contracts::ObjectOperations>::ObjectOptions;
@@ -105,15 +100,11 @@ where
// write has to fail with to abort the transaction.
let mut rejected = None;
// `update_config_with` refuses to run this closure over an unparseable
// tag set: merging onto it would silently drop every tag the bucket has —
// including the container ACL and versioning tags — because the rewrite
// closures treat "no parsed tags" as "no tags".
let result = update_config_with(&bucket, BUCKET_TAGGING_CONFIG, |bm| {
// Merging onto an unparseable tag set would silently drop every tag
// the bucket has — including the container ACL and versioning tags —
// because the rewrite closures treat "no parsed tags" as "no tags".
// Refuse instead: the persisted config is intact, just unreadable.
if !bm.tagging_config_xml.is_empty() && bm.tagging_config.is_none() {
return Err(SwiftStorageError::other(UNREADABLE_TAGGING_SENTINEL));
}
let tagging = match rewrite(bm.tagging_config.as_ref()) {
Ok(tagging) => tagging,
Err(err) => {
@@ -140,7 +131,7 @@ where
}
if let Err(err) = result {
let unreadable = err.to_string().contains(UNREADABLE_TAGGING_SENTINEL);
let unreadable = is_unreadable_config_error(&err);
tracing::error!(
event = EVENT_SWIFT_BUCKET_TAGGING_UPDATE,
component = LOG_COMPONENT_PROTOCOLS,
@@ -19,6 +19,7 @@
- `backlog-2102` rc.2/rc.3 empty scanner usage floor recovery: old DeleteBucket cleanup could synthesize an empty incomplete v2 usage primary/backup before leadership added an epoch, while newer scanners require a durable authoritative baseline identity. New scanners recognize only that exact serialized empty-fence shape, preserve its epoch through a CAS-protected recovery marker, and rebuild namespace coverage without treating zero usage as authoritative. Remove this recovery path and marker after rc.2 and rc.3 are no longer supported direct-upgrade sources.
- `backlog-2122` rc.1-rc.3 non-empty scanner usage floor recovery: leadership fencing in those releases can stamp scanner_epoch onto a real bucket-usage snapshot before any scanner cycle completed, leaving a non-empty floor with no scanner_cycle and no authoritative baseline identity. New scanners recognize only this consistent incomplete fenced shape, preserve the epoch through the CAS-protected recovery marker, and rebuild namespace coverage without treating the old usage data as authoritative. Remove this recovery path after rc.1, rc.2, and rc.3 are no longer supported direct-upgrade sources.
- `s3gate-metadata-xml` persisted bucket XML migration: mixed-version site-replication peers, retained `.metadata.bin` objects, and backup archives can all carry XML written by the s3s codec, so the gateway migration must keep the legacy codec available until every stored form has crossed a verified rewrite boundary. Remove the legacy s3s parser and serializer only after the minimum supported direct-upgrade release reads and writes every persisted XML configuration family through the gateway codec, every supported mixed-version site-replication topology has completed its writer upgrade, and migration tooling has verified or rewritten every retained bucket metadata object and restorable backup archive.
- `s3gate-parse-strict` unreadable bucket versioning on object writes: a stored versioning configuration that cannot be parsed was historically read as unversioned, and existing clusters may already hold such buckets unnoticed. RUSTFS_BUCKET_CONFIG_PARSE_MODE defaults to permissive, which keeps that unversioned write and records it through rustfs_bucket_metadata_parse_failed_total, rustfs_bucket_metadata_unparsable_current, and an error log; strict refuses the write with 503. Paths that already refuse an unreadable configuration (typed getters, the read-modify-write guard, Object Lock, default encryption, delete-time versioning) refuse in both modes. Mixed-version risk window: older nodes lack the per-config read-modify-write guard and still write unversioned, so the metric must be driven to zero before relying on strict. Flip the default to strict after the parse-failure counter has read zero fleet-wide for two consecutive releases, and remove the permissive branch and the switch one release after that.
- `rustfs-6339` legacy bucket policy ID casing: earlier RustFS releases persisted the top-level policy identifier as "ID", while current writes use the S3-compatible "Id" spelling. Readers accept both spellings so retained bucket metadata remains usable after upgrade. Remove the legacy alias after migration tooling has rewritten every retained bucket policy using "ID".
- `table-publication-fence-v1` table publication fencing: nodes that predate table and table-bucket publication fences can mutate live files while a new node is publishing a catalog pointer. New nodes retain exact object guards until the operator confirms that every serving node uses the new fences. Fleet confirmation also requires non-overlapping active warehouse prefixes and lifecycle workers that exclude table buckets. Remove the exact live-file fallback and the fleet-confirmation gate after the minimum supported RustFS release acquires table fences for registered-table mutations and table-bucket fences for unresolved-prefix mutations.
- `table-catalog-strong-snapshot-v1` durable strong catalog snapshot compatibility: version 1 writes continue during mixed-version rollout until operators confirm that every serving node reads version 2, and version 1 table/view identifier collisions remain available only for cleanup. Remove version 1 writes and collision cleanup after the minimum supported RustFS release reads version 2 and every retained durable strong snapshot is collision-free and has been upgraded to version 2.
+21 -12
View File
@@ -1914,9 +1914,11 @@ impl DefaultBucketUsecase {
let rules = match metadata_sys::get_lifecycle_config(&bucket).await {
Ok((cfg, _)) => cfg.rules,
Err(_) => {
Err(StorageError::ConfigNotFound) => {
return Err(s3_error!(NoSuchLifecycleConfiguration));
}
// An unreadable stored configuration is not an absent one.
Err(err) => return Err(ApiError::from(err).into()),
};
Ok(S3Response::new(GetBucketLifecycleConfigurationOutput {
@@ -1940,17 +1942,24 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
let has_notification_config = metadata_sys::get_notification_config(&bucket).await.unwrap_or_else(|err| {
warn!(
component = LOG_COMPONENT_APP,
subsystem = LOG_SUBSYSTEM_BUCKET,
event = "bucket_notification_config_load_failed",
bucket = %bucket,
error = ?err,
"Failed to load bucket notification configuration"
);
None
});
let has_notification_config = match metadata_sys::get_notification_config(&bucket).await {
Ok(config) => config,
// An unreadable config is not "no notifications configured".
Err(err) if crate::storage_api::error::is_unreadable_config_error(&err) => {
return Err(ApiError::from(err).into());
}
Err(err) => {
warn!(
component = LOG_COMPONENT_APP,
subsystem = LOG_SUBSYSTEM_BUCKET,
event = "bucket_notification_config_load_failed",
bucket = %bucket,
error = ?err,
"Failed to load bucket notification configuration"
);
None
}
};
if let Some(NotificationConfiguration {
event_bridge_configuration,
@@ -3155,7 +3155,7 @@ async fn delete_object_versioning_config_failure_leaves_latest_object_intact() {
.await
.expect_err("versioning config failure must reject DeleteObject");
assert_eq!(err.code(), &s3s::S3ErrorCode::InternalError);
assert_eq!(err.code(), &s3s::S3ErrorCode::ServiceUnavailable);
assert_eq!(read_object_bytes(&ecstore, &bucket, object).await, payload);
}
+3 -2
View File
@@ -629,8 +629,9 @@ impl DefaultMultipartUsecase {
let mut opts = get_complete_multipart_upload_opts_with_replication_authorization(&req.headers, replication_authorized)
.map_err(ApiError::from)?;
apply_bucket_generation_guard(&req, &bucket, &mut opts)?;
let versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
let version_suspended = BucketVersioningSys::prefix_suspended(&bucket, &key).await;
let (versioned, version_suspended) = BucketVersioningSys::write_state(&bucket, &key)
.await
.map_err(ApiError::from)?;
opts.versioned = versioned;
opts.version_suspended = version_suspended;
let capacity_scope_token = Uuid::new_v4();
+1 -1
View File
@@ -1429,7 +1429,7 @@ mod tests {
.await
.expect_err("an unreadable bucket encryption configuration must refuse the copy");
assert_eq!(err.code(), &S3ErrorCode::InternalError);
assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable);
let lookup_err = store
.get_object_info(&bucket, destination, &ObjectOptions::default())
.await
+4 -2
View File
@@ -554,9 +554,11 @@ impl DefaultObjectUsecase {
opts.expected_bucket_incarnation_id = ctx.expected_bucket_incarnation_id;
opts.preserve_etag = ctx.preserve_etag.clone();
opts.preserve_delete_marker = ctx.preserve_delete_marker;
let versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
let (versioned, version_suspended) = BucketVersioningSys::write_state(&bucket, &key)
.await
.map_err(ApiError::from)?;
opts.versioned = versioned;
opts.version_suspended = BucketVersioningSys::prefix_suspended(&bucket, &key).await;
opts.version_suspended = version_suspended;
let capacity_scope_token = Uuid::new_v4();
opts.capacity_scope_token = Some(capacity_scope_token);
+2 -2
View File
@@ -4187,7 +4187,7 @@ mod tests {
.await
.expect_err("an unreadable bucket encryption configuration must refuse the write");
assert_eq!(err.code(), &S3ErrorCode::InternalError);
assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable);
let lookup_err = store
.get_object_info(&bucket, object, &ObjectOptions::default())
.await
@@ -4268,7 +4268,7 @@ mod tests {
.await
.expect_err("an unreadable bucket encryption configuration must refuse the extract upload");
assert_eq!(err.code(), &S3ErrorCode::InternalError);
assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable);
let lookup_err = store
.get_object_info(&bucket, "archive.tar", &ObjectOptions::default())
.await
+4 -3
View File
@@ -384,8 +384,9 @@ pub(super) fn resolve_bucket_default_sse(
/// bucket whose metadata document is absent — `ConfigNotFound`, so a cold
/// cache and a missing bucket are never turned into a refusal, and the write
/// still fails later with its own `NoSuchBucket`;
/// * blob present but unparseable — deterministic, so retrying cannot help;
/// surfaces as `InternalError` until an operator repairs or removes it;
/// * blob present but unparseable — the typed unreadable-config refusal,
/// surfaced as `ServiceUnavailable` naming the bucket and config until an
/// operator repairs or removes it (rustfs/backlog#1734);
/// * the metadata read itself failed (namespace lock, quorum, disk, an
/// uninitialized metadata system) — transient, and the typed error maps to
/// the retryable `ServiceUnavailable`.
@@ -955,7 +956,7 @@ impl DefaultObjectUsecase {
pub(crate) async fn object_lock_checks_required(bucket: &str) -> bool {
get_bucket_metadata(bucket)
.await
.map_or(true, |metadata| metadata.object_locking())
.map_or(true, |metadata| metadata.object_lock_checks_required())
}
pub(super) fn object_lock_checks_required_for_state(state: &metadata_sys::ObjectLockConfigState) -> bool {
+30 -1
View File
@@ -13,7 +13,7 @@
// limitations under the License.
use crate::storage_api::error::contract::{StorageErrorCode, range::HTTPRangeError};
use crate::storage_api::error::{PoolMetadataError, QuotaError, StorageError};
use crate::storage_api::error::{PoolMetadataError, QuotaError, StorageError, unreadable_config_refusal};
use http::StatusCode;
use rustfs_kms::KmsUnavailableError;
use s3s::{S3Error, S3ErrorCode};
@@ -532,6 +532,16 @@ impl From<StorageError> for ApiError {
source: Some(Box::new(err)),
};
}
// A stored bucket config that cannot be parsed stays refused until an
// operator repairs it, so it is a recoverable 503 that names what to
// repair rather than a generic 500 (rustfs/backlog#1734).
if let Some(message) = unreadable_config_refusal(&err).map(ToString::to_string) {
return ApiError {
code: S3ErrorCode::ServiceUnavailable,
message,
source: Some(Box::new(err)),
};
}
if let StorageError::Io(ref io_err) = err
&& let Some(inner) = io_err.get_ref()
{
@@ -817,6 +827,25 @@ mod tests {
use s3s::{S3Error, S3ErrorCode};
use std::io::{Error as IoError, ErrorKind};
/// rustfs/backlog#1734: a refusal to act on a stored bucket config that
/// cannot be parsed is recoverable once an operator repairs the bytes, so
/// it answers 503 naming the bucket, config and stored length, not a
/// generic 500. It must survive the error being cloned on the way.
#[test]
fn unreadable_bucket_config_maps_to_service_unavailable_naming_the_config() {
let err = StorageError::other(crate::storage_api::error::UnreadableBucketConfig {
bucket: "photos".to_string(),
config_file: "versioning.xml".to_string(),
raw_len: 42,
});
for api_error in [ApiError::from(err.clone()), ApiError::from(err)] {
assert_eq!(api_error.code, S3ErrorCode::ServiceUnavailable);
for needle in ["photos", "versioning.xml", "42"] {
assert!(api_error.message.contains(needle), "{needle} missing from {:?}", api_error.message);
}
}
}
#[test]
fn api_error_diagnostic_preserves_typed_cause_without_sensitive_payload() {
let error = ApiError::from(StorageError::Io(IoError::new(ErrorKind::TimedOut, "secret=do-not-log")));
+71 -9
View File
@@ -15,6 +15,7 @@
use crate::runtime_sources::current_region;
use crate::server::ShutdownHandle;
use crate::server::runtime_sources::current_notify_interface;
use crate::storage_api::error::{StorageError, is_unreadable_config_error};
use crate::storage_api::startup::bucket_metadata::contract::bucket::{BucketOperations, BucketOptions};
use crate::storage_api::startup::init::{
get_bucket_notification_config, process_lambda_configurations, process_queue_configurations, process_topic_configurations,
@@ -169,13 +170,48 @@ fn notification_config_to_event_rules(
Ok(event_rules)
}
async fn apply_bucket_notification_configuration(bucket: &str, region: &str) -> Result<bool, NotificationError> {
let has_notification_config = get_bucket_notification_config(bucket)
.await
.map_err(|err| NotificationError::StorageNotAvailable(format!("load bucket notification config for {bucket}: {err}")))?;
/// One bucket's persisted notification configuration as startup sees it.
#[derive(Debug)]
enum BucketNotificationLookup {
Configured(s3s::dto::NotificationConfiguration),
Missing,
/// Stored bytes cannot be parsed. Deterministic, so a retry cannot help,
/// and reading it as missing would clear the bucket's rules.
Unreadable,
}
match has_notification_config {
Some(cfg) => {
fn classify_bucket_notification_lookup(
bucket: &str,
lookup: Result<Option<s3s::dto::NotificationConfiguration>, StorageError>,
) -> Result<BucketNotificationLookup, NotificationError> {
match lookup {
Ok(Some(cfg)) => Ok(BucketNotificationLookup::Configured(cfg)),
Ok(None) => Ok(BucketNotificationLookup::Missing),
Err(err) if is_unreadable_config_error(&err) => {
error!(
target: "rustfs::init",
event = "notification_config_unreadable",
component = LOG_COMPONENT_INIT,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
bucket = %bucket,
error = %err,
"Bucket notification configuration is unreadable; leaving its rules unchanged"
);
Ok(BucketNotificationLookup::Unreadable)
}
Err(err) => Err(NotificationError::StorageNotAvailable(format!(
"load bucket notification config for {bucket}: {err}"
))),
}
}
/// Apply one bucket's persisted notification rules. An unreadable config is
/// isolated to its bucket: it neither aborts setup for the other buckets nor
/// clears the bucket's rules as if it had none.
async fn apply_bucket_notification_configuration(bucket: &str, region: &str) -> Result<bool, NotificationError> {
match classify_bucket_notification_lookup(bucket, get_bucket_notification_config(bucket).await)? {
BucketNotificationLookup::Unreadable => Ok(false),
BucketNotificationLookup::Configured(cfg) => {
info!(
target: "rustfs::init",
event = "notification_config_loaded",
@@ -195,7 +231,7 @@ async fn apply_bucket_notification_configuration(bucket: &str, region: &str) ->
.await?;
Ok(true)
}
None => {
BucketNotificationLookup::Missing => {
info!(
target: "rustfs::init",
event = "notification_config_missing",
@@ -1375,16 +1411,42 @@ pub async fn init_sftp_system() -> Result<Option<ShutdownHandle>, Box<dyn std::e
#[cfg(test)]
mod tests {
use super::{
build_aws_kms_config, build_vault_kms_config, build_vault_transit_kms_config, notification_config_to_event_rules,
resolve_buffer_profile_config,
BucketNotificationLookup, build_aws_kms_config, build_vault_kms_config, build_vault_transit_kms_config,
classify_bucket_notification_lookup, notification_config_to_event_rules, resolve_buffer_profile_config,
};
use crate::config::{BufferConfig, WorkloadProfile};
use crate::storage_api::error::{StorageError, UnreadableBucketConfig};
use rustfs_config::KI_B;
use rustfs_s3_types::EventName;
use s3s::dto::{
FilterRule, FilterRuleName, NotificationConfiguration, NotificationConfigurationFilter, QueueConfiguration, S3KeyFilter,
};
/// rustfs/backlog#1734: an unreadable notification config is isolated to
/// its bucket. It must not abort setup for every other bucket (a storage
/// fault still does, so reconciliation retries) and must not read as
/// missing, which would clear the bucket's rules.
#[test]
fn unreadable_notification_config_is_isolated_to_its_bucket() {
let unreadable = StorageError::other(UnreadableBucketConfig {
bucket: "b".to_string(),
config_file: "notification.xml".to_string(),
raw_len: 5,
});
assert!(matches!(
classify_bucket_notification_lookup("b", Err(unreadable)),
Ok(BucketNotificationLookup::Unreadable)
));
assert!(matches!(
classify_bucket_notification_lookup("b", Ok(None)),
Ok(BucketNotificationLookup::Missing)
));
assert!(matches!(
classify_bucket_notification_lookup("b", Err(StorageError::ErasureReadQuorum)),
Err(rustfs_notify::NotificationError::StorageNotAvailable(_))
));
}
#[test]
fn resolve_buffer_profile_config_returns_fallback_when_primary_is_invalid() {
let invalid_primary = WorkloadProfile::Custom(BufferConfig {
+86 -1
View File
@@ -21,6 +21,7 @@ use crate::storage_api::startup::bucket_metadata::{
reconcile_bucket_resync_target_intents, try_migrate_bucket_metadata, try_migrate_iam_config,
};
use std::{
future::Future,
io::{Error as IoError, Result as IoResult},
sync::Arc,
time::{Duration, Instant},
@@ -35,6 +36,8 @@ const EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_STARTED: &str = "replication_r
const LOG_COMPONENT_STARTUP_BUCKET_METADATA: &str = "startup_bucket_metadata";
const LOG_SUBSYSTEM_ON_DEMAND_MIGRATION: &str = "on_demand_migration";
const LOG_SUBSYSTEM_REPLICATION: &str = "replication";
const IAM_MIGRATION_MAX_RETRIES: usize = 15;
const IAM_MIGRATION_RETRY_INTERVAL: Duration = Duration::from_secs(1);
const METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_DURATION_SECONDS: &str =
"rustfs_replication_resync_startup_background_duration_seconds";
const METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_EVENTS_TOTAL: &str =
@@ -84,7 +87,7 @@ pub(crate) async fn init_bucket_metadata_runtime(store: Arc<ECStore>, ctx: Cance
try_migrate_bucket_metadata(store.clone()).await?;
try_migrate_iam_config(store.clone()).await?;
retry_iam_config_migration(|| try_migrate_iam_config(store.clone())).await?;
init_on_demand_migration_runtime();
init_bucket_metadata_sys(store, buckets.clone()).await;
spawn_bucket_resync_startup_reconcile(buckets.clone(), ctx, true);
@@ -92,6 +95,51 @@ pub(crate) async fn init_bucket_metadata_runtime(store: Arc<ECStore>, ctx: Cance
Ok(buckets)
}
async fn retry_iam_config_migration<Operation, OperationFuture>(mut operation: Operation) -> IoResult<()>
where
Operation: FnMut() -> OperationFuture,
OperationFuture: Future<Output = IoResult<()>>,
{
retry_iam_config_migration_with(&mut operation, IAM_MIGRATION_MAX_RETRIES, IAM_MIGRATION_RETRY_INTERVAL).await
}
async fn retry_iam_config_migration_with<Operation, OperationFuture>(
operation: &mut Operation,
max_retries: usize,
retry_interval: Duration,
) -> IoResult<()>
where
Operation: FnMut() -> OperationFuture,
OperationFuture: Future<Output = IoResult<()>>,
{
let mut retries = 0;
loop {
match operation().await {
Ok(()) => return Ok(()),
Err(error) if iam_migration_error_is_retryable(&error) && retries < max_retries => {
retries += 1;
tracing::warn!(
component = LOG_COMPONENT_STARTUP_BUCKET_METADATA,
subsystem = "iam_migration",
retry_count = retries,
max_retries,
error = %error,
"IAM config migration hit a transient quorum error; retrying"
);
tokio::time::sleep(retry_interval).await;
}
Err(error) => return Err(error),
}
}
}
fn iam_migration_error_is_retryable(error: &IoError) -> bool {
error
.get_ref()
.and_then(|source| source.downcast_ref::<StorageError>())
.is_some_and(StorageError::is_quorum_error)
}
/// Publishes the on-demand migration module switch, installs the app-layer
/// write-back the pull pipeline stores objects with (rustfs/backlog#2153),
/// and registers the runtime's config hook before bucket metadata is
@@ -281,4 +329,41 @@ mod tests {
assert_eq!(STARTUP_BACKGROUND_STATUS_RUNNING, 2.0);
assert_eq!(STARTUP_BACKGROUND_STATUS_CANCELED, 3.0);
}
#[tokio::test]
async fn iam_migration_retries_only_quorum_errors() {
let mut attempts = 0;
retry_iam_config_migration_with(
&mut || {
attempts += 1;
std::future::ready(if attempts == 1 {
Err(IoError::other(StorageError::InsufficientReadQuorum(
".minio.sys".into(),
"config/iam/".into(),
)))
} else {
Ok(())
})
},
1,
Duration::ZERO,
)
.await
.expect("quorum recovery must allow startup to continue");
assert_eq!(attempts, 2);
let mut deterministic_attempts = 0;
let error = retry_iam_config_migration_with(
&mut || {
deterministic_attempts += 1;
std::future::ready(Err(IoError::other("incompatible IAM metadata")))
},
1,
Duration::ZERO,
)
.await
.expect_err("deterministic migration errors must fail immediately");
assert_eq!(error.to_string(), "incompatible IAM metadata");
assert_eq!(deterministic_attempts, 1);
}
}
+10
View File
@@ -26,6 +26,7 @@ const EVENT_EXTERNAL_ENV_COMPAT_CONFLICT: &str = "external_env_compat_conflict";
const EVENT_EXTERNAL_ENV_COMPAT_APPLIED: &str = "external_env_compat_applied";
const EVENT_OBSERVABILITY_GUARD_SET: &str = "observability_guard_set";
const EVENT_OBSERVABILITY_GUARD_SET_FAILED: &str = "observability_guard_set_failed";
const EVENT_BUCKET_CONFIG_PARSE_MODE_SELECTED: &str = "bucket_config_parse_mode_selected";
#[derive(Debug)]
pub(crate) enum StartupServerPreflightError {
@@ -66,6 +67,15 @@ pub(crate) async fn init_startup_server_preflight(
init_license(config.license.clone());
init_startup_observability(config.obs_endpoint.clone()).await?;
log_external_prefix_compat_report(env_compat_report);
let parse_mode = crate::storage_api::startup::init::validate_bucket_config_parse_mode_env()
.map_err(|err| StartupServerPreflightError::Other(Error::other(err)))?;
info!(
event = EVENT_BUCKET_CONFIG_PARSE_MODE_SELECTED,
component = LOG_COMPONENT_MAIN,
subsystem = LOG_SUBSYSTEM_STARTUP,
mode = parse_mode.as_str(),
"Bucket config parse mode selected"
);
init_startup_runtime_foundation(config)
.await
.map_err(StartupServerPreflightError::Other)
+13 -2
View File
@@ -145,6 +145,17 @@ fn get_skip_verify_bitrot() -> bool {
}
}
/// Like [`bucket_versioning_config`], for laying out an object write or
/// delete: an unreadable stored configuration is refused in strict mode
/// instead of being read as unversioned (rustfs/backlog#1734).
async fn bucket_versioning_config_for_write(bucket: &str) -> Result<VersioningConfiguration> {
#[cfg(test)]
VERSIONING_CONFIG_LOOKUPS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
#[cfg(test)]
wait_for_versioning_config_test_hook(bucket).await;
BucketVersioningSys::get_for_write(bucket).await
}
/// Creates options for deleting an object in a bucket.
pub async fn del_opts(
bucket: &str,
@@ -153,7 +164,7 @@ pub async fn del_opts(
headers: &HeaderMap<HeaderValue>,
metadata: HashMap<String, String>,
) -> Result<ObjectOptions> {
let versioning_cfg = bucket_versioning_config(bucket).await;
let versioning_cfg = bucket_versioning_config_for_write(bucket).await?;
del_opts_with_versioning(bucket, object, vid, headers, metadata, &versioning_cfg, false)
}
@@ -350,7 +361,7 @@ pub async fn put_opts_with_replication_authorization(
metadata: HashMap<String, String>,
replication_request_authorized: bool,
) -> Result<ObjectOptions> {
let versioning_cfg = bucket_versioning_config(bucket).await;
let versioning_cfg = bucket_versioning_config_for_write(bucket).await?;
let versioned = versioning_cfg.prefix_enabled(object);
let version_suspended = versioning_cfg.prefix_suspended(object);
+1 -1
View File
@@ -414,7 +414,7 @@ pub(crate) mod ecstore_bucket {
bandwidth, bucket_target_sys, durability, lifecycle, metadata, metadata_sys, migration, object_lock, policy_sys,
remote_s3_client, replication, tagging, target, utils,
};
pub(crate) use rustfs_ecstore::api::bucket::{quota, versioning, versioning_sys};
pub(crate) use rustfs_ecstore::api::bucket::{config_parse_mode, quota, versioning, versioning_sys};
}
pub(crate) mod ecstore_capacity {
+6
View File
@@ -83,6 +83,11 @@ pub(crate) mod error {
}
}
#[cfg(test)]
pub(crate) use crate::storage::storage_api::ecstore_bucket::metadata::UnreadableBucketConfig;
pub(crate) use crate::storage::storage_api::ecstore_bucket::metadata::{
is_unreadable_config_error, unreadable_config_refusal,
};
pub(crate) use crate::storage::storage_api::ecstore_error::PoolMetadataError;
#[cfg(test)]
pub(crate) use crate::storage::storage_api::ecstore_error::PoolMetadataFailure;
@@ -326,6 +331,7 @@ pub(crate) mod startup {
}
pub(crate) mod init {
pub(crate) use crate::storage::storage_api::ecstore_bucket::config_parse_mode::validate_bucket_config_parse_mode_env;
pub(crate) use crate::storage::storage_api::{
get_bucket_notification_config, process_lambda_configurations, process_queue_configurations,
process_topic_configurations,