Merge branch 'main' into fix/auto-repair-20260901-005614

This commit is contained in:
Zhengchao An
2026-09-01 03:43:03 +08:00
committed by GitHub
27 changed files with 3971 additions and 577 deletions
Generated
+50 -50
View File
@@ -3916,7 +3916,7 @@ checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555"
[[package]]
name = "e2e_test"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"anyhow",
"astral-tokio-tar",
@@ -9386,7 +9386,7 @@ dependencies = [
[[package]]
name = "rustfs"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"aes-gcm",
"anyhow",
@@ -9527,7 +9527,7 @@ dependencies = [
[[package]]
name = "rustfs-audit"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"const-str",
"futures",
@@ -9549,7 +9549,7 @@ dependencies = [
[[package]]
name = "rustfs-checksums"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"base64-simd",
"bytes",
@@ -9565,7 +9565,7 @@ dependencies = [
[[package]]
name = "rustfs-common"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"metrics",
@@ -9578,7 +9578,7 @@ dependencies = [
[[package]]
name = "rustfs-concurrency"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"insta",
@@ -9591,7 +9591,7 @@ dependencies = [
[[package]]
name = "rustfs-config"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"const-str",
"hotpath",
@@ -9601,7 +9601,7 @@ dependencies = [
[[package]]
name = "rustfs-credentials"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"base64-simd",
"hmac 0.13.0",
@@ -9615,7 +9615,7 @@ dependencies = [
[[package]]
name = "rustfs-crypto"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"aes-gcm",
"argon2",
@@ -9636,7 +9636,7 @@ dependencies = [
[[package]]
name = "rustfs-data-usage"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"rmp-serde",
@@ -9646,7 +9646,7 @@ dependencies = [
[[package]]
name = "rustfs-ecstore"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"async-channel",
@@ -9781,7 +9781,7 @@ dependencies = [
[[package]]
name = "rustfs-extension-schema"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"serde",
@@ -9791,7 +9791,7 @@ dependencies = [
[[package]]
name = "rustfs-filemeta"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"byteorder",
@@ -9818,7 +9818,7 @@ dependencies = [
[[package]]
name = "rustfs-heal"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"base64-simd",
@@ -9854,7 +9854,7 @@ dependencies = [
[[package]]
name = "rustfs-heal-contracts"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"serde",
"serde_json",
@@ -9864,7 +9864,7 @@ dependencies = [
[[package]]
name = "rustfs-iam"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"async-trait",
@@ -9912,7 +9912,7 @@ dependencies = [
[[package]]
name = "rustfs-io-core"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"bytes",
"hotpath",
@@ -9924,7 +9924,7 @@ dependencies = [
[[package]]
name = "rustfs-io-metrics"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"criterion",
"hotpath",
@@ -9988,7 +9988,7 @@ dependencies = [
[[package]]
name = "rustfs-keystone"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"bytes",
"futures",
@@ -10015,7 +10015,7 @@ dependencies = [
[[package]]
name = "rustfs-kms"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"aes-gcm",
"anyhow",
@@ -10065,7 +10065,7 @@ dependencies = [
[[package]]
name = "rustfs-lifecycle"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"hotpath",
@@ -10088,7 +10088,7 @@ dependencies = [
[[package]]
name = "rustfs-lock"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"compact_str",
@@ -10111,7 +10111,7 @@ dependencies = [
[[package]]
name = "rustfs-log-analyzer"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"chrono",
"flate2",
@@ -10130,7 +10130,7 @@ dependencies = [
[[package]]
name = "rustfs-madmin"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"http 1.5.0",
@@ -10168,7 +10168,7 @@ dependencies = [
[[package]]
name = "rustfs-notify"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"async-trait",
@@ -10203,7 +10203,7 @@ dependencies = [
[[package]]
name = "rustfs-object-capacity"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"criterion",
"futures",
@@ -10222,7 +10222,7 @@ dependencies = [
[[package]]
name = "rustfs-object-data-cache"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"bytes",
"criterion",
@@ -10239,7 +10239,7 @@ dependencies = [
[[package]]
name = "rustfs-obs"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"chrono",
"crossbeam-channel",
@@ -10297,7 +10297,7 @@ dependencies = [
[[package]]
name = "rustfs-policy"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"base64-simd",
@@ -10328,7 +10328,7 @@ dependencies = [
[[package]]
name = "rustfs-protocols"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"astral-tokio-tar",
"async-compression",
@@ -10390,7 +10390,7 @@ dependencies = [
[[package]]
name = "rustfs-protos"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"flatbuffers",
"hotpath",
@@ -10415,7 +10415,7 @@ dependencies = [
[[package]]
name = "rustfs-replication"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"byteorder",
"bytes",
@@ -10433,7 +10433,7 @@ dependencies = [
[[package]]
name = "rustfs-rio"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"aes-gcm",
"arc-swap",
@@ -10474,7 +10474,7 @@ dependencies = [
[[package]]
name = "rustfs-rio-v2"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"aes-gcm",
"bytes",
@@ -10497,7 +10497,7 @@ dependencies = [
[[package]]
name = "rustfs-s3-client"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"base64-simd",
"bytes",
@@ -10541,7 +10541,7 @@ dependencies = [
[[package]]
name = "rustfs-s3-ops"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"rustfs-s3-types",
@@ -10549,7 +10549,7 @@ dependencies = [
[[package]]
name = "rustfs-s3-types"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"serde",
@@ -10558,7 +10558,7 @@ dependencies = [
[[package]]
name = "rustfs-s3select-api"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"async-compression",
@@ -10593,7 +10593,7 @@ dependencies = [
[[package]]
name = "rustfs-s3select-query"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-recursion",
"async-trait",
@@ -10612,7 +10612,7 @@ dependencies = [
[[package]]
name = "rustfs-scanner"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"bytes",
@@ -10655,7 +10655,7 @@ dependencies = [
[[package]]
name = "rustfs-scanner-contracts"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"chrono",
"jiff",
@@ -10670,7 +10670,7 @@ dependencies = [
[[package]]
name = "rustfs-security-governance"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"thiserror 2.0.20",
@@ -10678,7 +10678,7 @@ dependencies = [
[[package]]
name = "rustfs-signer"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"base64-simd",
"bytes",
@@ -10696,7 +10696,7 @@ dependencies = [
[[package]]
name = "rustfs-storage-api"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"hotpath",
@@ -10711,7 +10711,7 @@ dependencies = [
[[package]]
name = "rustfs-targets"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"async-nats",
@@ -10765,7 +10765,7 @@ dependencies = [
[[package]]
name = "rustfs-test-utils"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"hotpath",
"rustfs-data-usage",
@@ -10781,7 +10781,7 @@ dependencies = [
[[package]]
name = "rustfs-tls-runtime"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"arc-swap",
"hotpath",
@@ -10802,7 +10802,7 @@ dependencies = [
[[package]]
name = "rustfs-trusted-proxies"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-trait",
"axum",
@@ -10839,7 +10839,7 @@ dependencies = [
[[package]]
name = "rustfs-utils"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"base64-simd",
"blake2",
@@ -10881,7 +10881,7 @@ dependencies = [
[[package]]
name = "rustfs-zip"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
dependencies = [
"async-compression",
"hotpath",
+50 -50
View File
@@ -72,7 +72,7 @@ edition = "2024"
license = "Apache-2.0"
repository = "https://github.com/rustfs/rustfs"
rust-version = "1.97.1"
version = "1.0.0-rc.4"
version = "1.0.0-rc.5"
homepage = "https://rustfs.com"
description = "RustFS is a high-performance distributed object storage software built using Rust, one of the most popular languages worldwide. "
keywords = ["RustFS", "Minio", "object-storage", "filesystem", "s3"]
@@ -89,55 +89,55 @@ redundant_clone = "warn"
[workspace.dependencies]
# RustFS Internal Crates
rustfs = { path = "./rustfs", version = "1.0.0-rc.4" }
rustfs-heal = { path = "crates/heal", version = "1.0.0-rc.4" }
rustfs-heal-contracts = { path = "crates/heal-contracts", version = "1.0.0-rc.4" }
rustfs-scanner-contracts = { path = "crates/scanner-contracts", version = "1.0.0-rc.4" }
rustfs-audit = { path = "crates/audit", version = "1.0.0-rc.4" }
rustfs-checksums = { path = "crates/checksums", version = "1.0.0-rc.4" }
rustfs-common = { path = "crates/common", version = "1.0.0-rc.4" }
rustfs-data-usage = { path = "crates/data-usage", version = "1.0.0-rc.4" }
rustfs-config = { path = "./crates/config", version = "1.0.0-rc.4" }
rustfs-concurrency = { path = "./crates/concurrency", version = "1.0.0-rc.4" }
rustfs-credentials = { path = "crates/credentials", version = "1.0.0-rc.4" }
rustfs-crypto = { path = "crates/crypto", version = "1.0.0-rc.4" }
rustfs-ecstore = { path = "crates/ecstore", version = "1.0.0-rc.4" }
rustfs-filemeta = { path = "crates/filemeta", version = "1.0.0-rc.4" }
rustfs-iam = { path = "crates/iam", version = "1.0.0-rc.4" }
rustfs-keystone = { path = "crates/keystone", version = "1.0.0-rc.4" }
rustfs-lifecycle = { path = "crates/lifecycle", version = "1.0.0-rc.4" }
rustfs-kms = { path = "crates/kms", version = "1.0.0-rc.4" }
rustfs-lock = { path = "crates/lock", version = "1.0.0-rc.4" }
rustfs-madmin = { path = "crates/madmin", version = "1.0.0-rc.4" }
rustfs-notify = { path = "crates/notify", version = "1.0.0-rc.4" }
rustfs-io-metrics = { path = "crates/io-metrics", version = "1.0.0-rc.4" }
rustfs-io-core = { path = "crates/io-core", version = "1.0.0-rc.4" }
rustfs-object-capacity = { path = "crates/object-capacity", version = "1.0.0-rc.4" }
rustfs-object-data-cache = { path = "crates/object-data-cache", version = "1.0.0-rc.4", default-features = false }
rustfs-log-analyzer = { path = "crates/log-analyzer", version = "1.0.0-rc.4" }
rustfs-obs = { path = "crates/obs", version = "1.0.0-rc.4" }
rustfs-policy = { path = "crates/policy", version = "1.0.0-rc.4" }
rustfs-protos = { path = "crates/protos", version = "1.0.0-rc.4" }
rustfs-protocols = { path = "crates/protocols", version = "1.0.0-rc.4" }
rustfs-replication = { path = "crates/replication", version = "1.0.0-rc.4" }
rustfs-rio = { path = "crates/rio", version = "1.0.0-rc.4" }
rustfs-rio-v2 = { path = "crates/rio-v2", version = "1.0.0-rc.4" }
rustfs-s3-client = { path = "crates/s3-client", version = "1.0.0-rc.4" }
rustfs-s3-types = { path = "crates/s3-types", version = "1.0.0-rc.4" }
rustfs-s3-ops = { path = "crates/s3-ops", version = "1.0.0-rc.4" }
rustfs-s3select-api = { path = "crates/s3select-api", version = "1.0.0-rc.4" }
rustfs-s3select-query = { path = "crates/s3select-query", version = "1.0.0-rc.4" }
rustfs-scanner = { path = "crates/scanner", version = "1.0.0-rc.4" }
rustfs-security-governance = { path = "crates/security-governance", version = "1.0.0-rc.4" }
rustfs-extension-schema = { path = "crates/extension-schema", version = "1.0.0-rc.4" }
rustfs-signer = { path = "crates/signer", version = "1.0.0-rc.4" }
rustfs-storage-api = { path = "crates/storage-api", version = "1.0.0-rc.4" }
rustfs-trusted-proxies = { path = "crates/trusted-proxies", version = "1.0.0-rc.4" }
rustfs-targets = { path = "crates/targets", version = "1.0.0-rc.4" }
rustfs-test-utils = { path = "crates/test-utils", version = "1.0.0-rc.4" }
rustfs-tls-runtime = { path = "crates/tls-runtime", version = "1.0.0-rc.4" }
rustfs-utils = { path = "crates/utils", version = "1.0.0-rc.4" }
rustfs-zip = { path = "./crates/zip", version = "1.0.0-rc.4" }
rustfs = { path = "./rustfs", version = "1.0.0-rc.5" }
rustfs-heal = { path = "crates/heal", version = "1.0.0-rc.5" }
rustfs-heal-contracts = { path = "crates/heal-contracts", version = "1.0.0-rc.5" }
rustfs-scanner-contracts = { path = "crates/scanner-contracts", version = "1.0.0-rc.5" }
rustfs-audit = { path = "crates/audit", version = "1.0.0-rc.5" }
rustfs-checksums = { path = "crates/checksums", version = "1.0.0-rc.5" }
rustfs-common = { path = "crates/common", version = "1.0.0-rc.5" }
rustfs-data-usage = { path = "crates/data-usage", version = "1.0.0-rc.5" }
rustfs-config = { path = "./crates/config", version = "1.0.0-rc.5" }
rustfs-concurrency = { path = "./crates/concurrency", version = "1.0.0-rc.5" }
rustfs-credentials = { path = "crates/credentials", version = "1.0.0-rc.5" }
rustfs-crypto = { path = "crates/crypto", version = "1.0.0-rc.5" }
rustfs-ecstore = { path = "crates/ecstore", version = "1.0.0-rc.5" }
rustfs-filemeta = { path = "crates/filemeta", version = "1.0.0-rc.5" }
rustfs-iam = { path = "crates/iam", version = "1.0.0-rc.5" }
rustfs-keystone = { path = "crates/keystone", version = "1.0.0-rc.5" }
rustfs-lifecycle = { path = "crates/lifecycle", version = "1.0.0-rc.5" }
rustfs-kms = { path = "crates/kms", version = "1.0.0-rc.5" }
rustfs-lock = { path = "crates/lock", version = "1.0.0-rc.5" }
rustfs-madmin = { path = "crates/madmin", version = "1.0.0-rc.5" }
rustfs-notify = { path = "crates/notify", version = "1.0.0-rc.5" }
rustfs-io-metrics = { path = "crates/io-metrics", version = "1.0.0-rc.5" }
rustfs-io-core = { path = "crates/io-core", version = "1.0.0-rc.5" }
rustfs-object-capacity = { path = "crates/object-capacity", version = "1.0.0-rc.5" }
rustfs-object-data-cache = { path = "crates/object-data-cache", version = "1.0.0-rc.5", default-features = false }
rustfs-log-analyzer = { path = "crates/log-analyzer", version = "1.0.0-rc.5" }
rustfs-obs = { path = "crates/obs", version = "1.0.0-rc.5" }
rustfs-policy = { path = "crates/policy", version = "1.0.0-rc.5" }
rustfs-protos = { path = "crates/protos", version = "1.0.0-rc.5" }
rustfs-protocols = { path = "crates/protocols", version = "1.0.0-rc.5" }
rustfs-replication = { path = "crates/replication", version = "1.0.0-rc.5" }
rustfs-rio = { path = "crates/rio", version = "1.0.0-rc.5" }
rustfs-rio-v2 = { path = "crates/rio-v2", version = "1.0.0-rc.5" }
rustfs-s3-client = { path = "crates/s3-client", version = "1.0.0-rc.5" }
rustfs-s3-types = { path = "crates/s3-types", version = "1.0.0-rc.5" }
rustfs-s3-ops = { path = "crates/s3-ops", version = "1.0.0-rc.5" }
rustfs-s3select-api = { path = "crates/s3select-api", version = "1.0.0-rc.5" }
rustfs-s3select-query = { path = "crates/s3select-query", version = "1.0.0-rc.5" }
rustfs-scanner = { path = "crates/scanner", version = "1.0.0-rc.5" }
rustfs-security-governance = { path = "crates/security-governance", version = "1.0.0-rc.5" }
rustfs-extension-schema = { path = "crates/extension-schema", version = "1.0.0-rc.5" }
rustfs-signer = { path = "crates/signer", version = "1.0.0-rc.5" }
rustfs-storage-api = { path = "crates/storage-api", version = "1.0.0-rc.5" }
rustfs-trusted-proxies = { path = "crates/trusted-proxies", version = "1.0.0-rc.5" }
rustfs-targets = { path = "crates/targets", version = "1.0.0-rc.5" }
rustfs-test-utils = { path = "crates/test-utils", version = "1.0.0-rc.5" }
rustfs-tls-runtime = { path = "crates/tls-runtime", version = "1.0.0-rc.5" }
rustfs-utils = { path = "crates/utils", version = "1.0.0-rc.5" }
rustfs-zip = { path = "./crates/zip", version = "1.0.0-rc.5" }
# Async Runtime and Networking
async-channel = "2.5.0"
+1 -1
View File
@@ -115,7 +115,7 @@ chown -R 10001:10001 data logs
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:latest
# Using specific version
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.4
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.5
```
If you use [podman](https://github.com/containers/podman) instead of docker, you can install the RustFS with the below command
+1 -1
View File
@@ -112,7 +112,7 @@ chown -R 10001:10001 data logs
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:latest
# 使用指定版本运行
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.4
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.5
```
如果您通过绑定挂载启用 TLS 证书目录,也请用同样方式准备该目录:
+135 -2
View File
@@ -4467,9 +4467,15 @@ async fn test_signed_put_object_extract_authorizes_each_pax_privilege_and_retent
let context_archive_resources = [
format!("arn:aws:s3:::{bucket}/tag-context.tar"),
format!("arn:aws:s3:::{bucket}/lock-context.tar"),
format!("arn:aws:s3:::{bucket}/legal-hold-context.tar"),
format!("arn:aws:s3:::{bucket}/user-agent-bypass.tar"),
format!("arn:aws:s3:::{bucket}/sse-bypass.tar"),
];
let tag_entry_resource = format!("arn:aws:s3:::{bucket}/tag-context-entry.txt");
let lock_entry_resource = format!("arn:aws:s3:::{bucket}/lock-context-entry.txt");
let legal_hold_entry_resource = format!("arn:aws:s3:::{bucket}/legal-hold-context-entry.txt");
let user_agent_entry_resource = format!("arn:aws:s3:::{bucket}/user-agent-bypass-entry.txt");
let sse_entry_resource = format!("arn:aws:s3:::{bucket}/sse-bypass-entry.txt");
let policy = serde_json::json!({
"Version": "2012-10-17",
"Statement": [
@@ -4529,7 +4535,7 @@ async fn test_signed_put_object_extract_authorizes_each_pax_privilege_and_retent
"Sid": "PaxContextArchives",
"Effect": "Allow",
"Principal": { "AWS": [pax_context_user] },
"Action": ["s3:PutObject", "s3:PutObjectRetention", "s3:PutObjectTagging"],
"Action": ["s3:PutObject", "s3:PutObjectRetention", "s3:PutObjectLegalHold", "s3:PutObjectTagging"],
"Resource": context_archive_resources
},
{
@@ -4569,6 +4575,49 @@ async fn test_signed_put_object_extract_authorizes_each_pax_privilege_and_retent
"Principal": { "AWS": [pax_context_user] },
"Action": ["s3:PutObjectRetention"],
"Resource": [lock_entry_resource]
},
{
"Sid": "PaxLegalHoldContextPut",
"Effect": "Allow",
"Principal": { "AWS": [pax_context_user] },
"Action": ["s3:PutObject"],
"Resource": [legal_hold_entry_resource.clone()]
},
{
"Sid": "PaxLegalHoldContextAction",
"Effect": "Allow",
"Principal": { "AWS": [pax_context_user] },
"Action": ["s3:PutObjectLegalHold"],
"Resource": [legal_hold_entry_resource],
"Condition": {
"StringEquals": {
"s3:object-lock-legal-hold": "OFF"
}
}
},
{
"Sid": "MemberUserAgentCondition",
"Effect": "Allow",
"Principal": { "AWS": [pax_context_user] },
"Action": ["s3:PutObject"],
"Resource": [user_agent_entry_resource],
"Condition": {
"StringEquals": {
"aws:UserAgent": "trusted"
}
}
},
{
"Sid": "MemberSseCondition",
"Effect": "Allow",
"Principal": { "AWS": [pax_context_user] },
"Action": ["s3:PutObject"],
"Resource": [sse_entry_resource],
"Condition": {
"StringEquals": {
"s3:x-amz-server-side-encryption": "AES256"
}
}
}
]
})
@@ -4581,9 +4630,14 @@ async fn test_signed_put_object_extract_authorizes_each_pax_privilege_and_retent
let cases = [
(
"legal-hold.tar",
put_only_client,
put_only_client.clone(),
HashMap::from([("minio.metadata.x-amz-object-lock-legal-hold", "ON".to_string())]),
),
(
"tagging.tar",
put_only_client,
HashMap::from([("minio.metadata.x-amz-tagging", "classification=restricted".to_string())]),
),
(
"retention-condition.tar",
conditional_client,
@@ -4670,6 +4724,57 @@ async fn test_signed_put_object_extract_authorizes_each_pax_privilege_and_retent
assert_eq!(stored.body.collect().await?.into_bytes().as_ref(), b"condition-body");
let pax_context_client = restricted_user_client(&env, pax_context_user, pax_context_secret);
for (archive_key, entry_key, pax_key, injected_value, outer_user_agent) in [
(
"user-agent-bypass.tar",
"user-agent-bypass-entry.txt",
"minio.metadata.user-agent",
"trusted",
Some("untrusted"),
),
(
"sse-bypass.tar",
"sse-bypass-entry.txt",
"minio.metadata.x-amz-server-side-encryption",
"AES256",
None,
),
] {
let pax = HashMap::from([(pax_key, injected_value.to_string())]);
let archive = make_tar_with_pax_entry(entry_key, b"must-not-write", None, &pax).await;
let err = pax_context_client
.put_object()
.bucket(bucket)
.key(archive_key)
.body(ByteStream::from(archive))
.customize()
.mutate_request(move |req| {
req.headers_mut().insert("x-amz-meta-snowball-auto-extract", "true");
if let Some(user_agent) = outer_user_agent {
req.headers_mut().insert("user-agent", user_agent);
}
})
.send()
.await
.expect_err("PAX metadata must not satisfy unrelated IAM request conditions");
assert_eq!(
err.as_service_error().and_then(|error| error.meta().code()),
Some("AccessDenied"),
"{archive_key}"
);
let err = admin_client
.head_object()
.bucket(bucket)
.key(entry_key)
.send()
.await
.expect_err("a denied PAX member must not be written");
assert!(matches!(
err.as_service_error().and_then(|error| error.meta().code()),
Some("NoSuchKey" | "NotFound")
));
}
let tag_pax = HashMap::from([("minio.metadata.x-amz-tagging", "classification=public".to_string())]);
let archive = make_tar_with_pax_entry("tag-context-entry.txt", b"tag-context-body", None, &tag_pax).await;
pax_context_client
@@ -4733,6 +4838,34 @@ async fn test_signed_put_object_extract_authorizes_each_pax_privilege_and_retent
pax_retain_until
);
let legal_hold_pax = HashMap::from([("minio.metadata.x-amz-object-lock-legal-hold", "ON".to_string())]);
let archive = make_tar_with_pax_entry("legal-hold-context-entry.txt", b"must-not-write", None, &legal_hold_pax).await;
let err = pax_context_client
.put_object()
.bucket(bucket)
.key("legal-hold-context.tar")
.object_lock_legal_hold_status(aws_sdk_s3::types::ObjectLockLegalHoldStatus::Off)
.body(ByteStream::from(archive))
.customize()
.mutate_request(|req| {
req.headers_mut().insert("x-amz-meta-snowball-auto-extract", "true");
})
.send()
.await
.expect_err("PAX legal hold must replace the outer value in the member IAM condition context");
assert_eq!(err.as_service_error().and_then(|error| error.meta().code()), Some("AccessDenied"));
let err = admin_client
.head_object()
.bucket(bucket)
.key("legal-hold-context-entry.txt")
.send()
.await
.expect_err("a denied PAX legal-hold member must not be written");
assert!(matches!(
err.as_service_error().and_then(|error| error.meta().code()),
Some("NoSuchKey" | "NotFound")
));
Ok(())
}
@@ -21,6 +21,53 @@ mod tests {
use std::error::Error;
use std::io::{Cursor, Write};
fn pax_record(key: &str, value: &str) -> Vec<u8> {
let payload = format!("{key}={value}\n");
let mut len = payload.len() + 3;
loop {
let record = format!("{len} {payload}");
if record.len() == len {
return record.into_bytes();
}
len = record.len();
}
}
async fn append_pax_header(
builder: &mut tokio_tar::Builder<Cursor<Vec<u8>>>,
entry_type: tokio_tar::EntryType,
records: &[(&str, &str)],
) -> Result<(), Box<dyn Error + Send + Sync>> {
let mut payload = Vec::new();
for (key, value) in records {
payload.extend(pax_record(key, value));
}
let mut header = tokio_tar::Header::new_ustar();
header.set_entry_type(entry_type);
header.set_size(u64::try_from(payload.len()).expect("PAX payload length should fit in u64"));
header.set_mode(0o644);
header.set_cksum();
builder
.append_data(&mut header, "PaxHeaders.X/snowball", Cursor::new(payload))
.await?;
Ok(())
}
async fn append_typed_entry(
builder: &mut tokio_tar::Builder<Cursor<Vec<u8>>>,
path: &str,
entry_type: tokio_tar::EntryType,
body: &[u8],
) -> Result<(), Box<dyn Error + Send + Sync>> {
let mut header = tokio_tar::Header::new_gnu();
header.set_entry_type(entry_type);
header.set_size(u64::try_from(body.len()).expect("TAR member length should fit in u64"));
header.set_mode(0o644);
header.set_cksum();
builder.append_data(&mut header, path, Cursor::new(body)).await?;
Ok(())
}
async fn build_test_archive() -> Result<Vec<u8>, Box<dyn Error + Send + Sync>> {
let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new()));
@@ -147,6 +194,57 @@ mod tests {
archive
}
async fn build_member_semantics_archive() -> Result<Vec<u8>, Box<dyn Error + Send + Sync>> {
let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new()));
append_pax_header(
&mut builder,
tokio_tar::EntryType::XGlobalHeader,
&[
("minio.metadata.x-amz-meta-owner", "global"),
("minio.metadata.x-amz-meta-snowball-auto-extract", "true"),
],
)
.await?;
append_pax_header(
&mut builder,
tokio_tar::EntryType::XHeader,
&[("minio.metadata.x-amz-meta-owner", "local")],
)
.await?;
append_typed_entry(&mut builder, "regular.txt", tokio_tar::EntryType::Regular, b"regular-body").await?;
for (path, entry_type) in [
("char", tokio_tar::EntryType::Char),
("block", tokio_tar::EntryType::Block),
("fifo", tokio_tar::EntryType::Fifo),
] {
append_typed_entry(&mut builder, path, entry_type, b"").await?;
}
let mut directory = tokio_tar::Header::new_gnu();
directory.set_entry_type(tokio_tar::EntryType::Directory);
directory.set_size(0);
directory.set_mode(0o755);
directory.set_cksum();
builder
.append_data(&mut directory, "directory/", Cursor::new(Vec::new()))
.await?;
for (path, entry_type) in [
("hard-link", tokio_tar::EntryType::Link),
("symlink", tokio_tar::EntryType::Symlink),
("continuous", tokio_tar::EntryType::Continuous),
("unknown", tokio_tar::EntryType::Other(b'9')),
] {
append_typed_entry(&mut builder, path, entry_type, b"").await?;
}
Ok(builder.into_inner().await?.into_inner())
}
async fn build_versioned_member_archive(path: &str, version_id: &str) -> Result<Vec<u8>, Box<dyn Error + Send + Sync>> {
let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new()));
append_pax_header(&mut builder, tokio_tar::EntryType::XHeader, &[("minio.versionId", version_id)]).await?;
append_typed_entry(&mut builder, path, tokio_tar::EntryType::Regular, b"versioned-body").await?;
Ok(builder.into_inner().await?.into_inner())
}
fn build_archive_with_invalid_utf8_entry() -> Vec<u8> {
let mut archive = Vec::new();
append_raw_tar_entry(&mut archive, b"invalid-\xff.txt", b"ignored-body");
@@ -199,6 +297,147 @@ mod tests {
Ok(())
}
#[tokio::test]
async fn snowball_auto_extract_applies_member_semantics_and_metadata_precedence() -> Result<(), Box<dyn Error + Send + Sync>>
{
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let client = env.create_s3_client();
let bucket = "snowball-member-semantics";
client.create_bucket().bucket(bucket).send().await?;
client
.put_object()
.bucket(bucket)
.key("fixture.tar")
.metadata("Snowball-Auto-Extract", "true")
.metadata("Minio-Snowball-Prefix", "members")
.metadata("owner", "outer")
.body(ByteStream::from(build_member_semantics_archive().await?))
.send()
.await?;
let regular = client.head_object().bucket(bucket).key("members/regular.txt").send().await?;
let regular_metadata = regular.metadata().expect("regular member should expose metadata");
assert_eq!(regular_metadata.get("owner").map(String::as_str), Some("local"));
assert!(!regular_metadata.contains_key("snowball-auto-extract"));
assert!(!regular_metadata.contains_key("minio-snowball-prefix"));
for key in ["char", "block", "fifo"] {
let head = client
.head_object()
.bucket(bucket)
.key(format!("members/{key}"))
.send()
.await?;
assert_eq!(head.content_length(), Some(0), "{key} should be materialized as an empty object");
assert_eq!(
head.metadata().and_then(|metadata| metadata.get("owner")).map(String::as_str),
Some("outer"),
"{key} should not inherit global PAX metadata"
);
}
let directory = client.head_object().bucket(bucket).key("members/directory/").send().await?;
assert_eq!(directory.content_length(), Some(0));
for key in ["hard-link", "symlink", "continuous", "unknown"] {
let error = client
.head_object()
.bucket(bucket)
.key(format!("members/{key}"))
.send()
.await
.expect_err("unsupported TAR entry type must be skipped");
assert_eq!(error.into_service_error().code(), Some("NotFound"), "{key}");
}
env.stop_server();
Ok(())
}
#[tokio::test]
async fn snowball_auto_extract_validates_pax_version_id_against_bucket_state() -> Result<(), Box<dyn Error + Send + Sync>> {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let client = env.create_s3_client();
let bucket = "snowball-version-semantics";
client.create_bucket().bucket(bucket).send().await?;
client
.put_object()
.bucket(bucket)
.key("null.tar")
.metadata("Snowball-Auto-Extract", "true")
.body(ByteStream::from(build_versioned_member_archive("null.txt", "null").await?))
.send()
.await?;
let null_member = client.get_object().bucket(bucket).key("null.txt").send().await?;
assert_eq!(null_member.body.collect().await?.into_bytes().as_ref(), b"versioned-body");
for (archive_key, member_key, version_id) in [
("uuid.tar", "uuid.txt", uuid::Uuid::new_v4().to_string()),
("uppercase-null.tar", "uppercase-null.txt", "NULL".to_string()),
] {
let error = client
.put_object()
.bucket(bucket)
.key(archive_key)
.metadata("Snowball-Auto-Extract", "true")
.body(ByteStream::from(build_versioned_member_archive(member_key, &version_id).await?))
.send()
.await
.expect_err("invalid or unversioned UUID import must be rejected");
assert_eq!(error.into_service_error().code(), Some("InvalidArgument"), "{archive_key}");
let missing = client
.head_object()
.bucket(bucket)
.key(member_key)
.send()
.await
.expect_err("rejected version import must not create an object");
assert_eq!(missing.into_service_error().code(), Some("NotFound"), "{member_key}");
}
client
.put_bucket_versioning()
.bucket(bucket)
.versioning_configuration(
aws_sdk_s3::types::VersioningConfiguration::builder()
.status(aws_sdk_s3::types::BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await?;
let imported_version_id = uuid::Uuid::new_v4().to_string();
client
.put_object()
.bucket(bucket)
.key("versioned-uuid.tar")
.metadata("Snowball-Auto-Extract", "true")
.body(ByteStream::from(
build_versioned_member_archive("versioned-uuid.txt", &imported_version_id).await?,
))
.send()
.await?;
let imported = client
.get_object()
.bucket(bucket)
.key("versioned-uuid.txt")
.version_id(&imported_version_id)
.send()
.await?;
assert_eq!(imported.version_id(), Some(imported_version_id.as_str()));
assert_eq!(imported.body.collect().await?.into_bytes().as_ref(), b"versioned-body");
env.stop_server();
Ok(())
}
#[tokio::test]
async fn snowball_auto_extract_supports_standard_headers_with_combined_extract_options()
-> Result<(), Box<dyn Error + Send + Sync>> {
@@ -198,6 +198,9 @@ pub enum S3KeyName {
#[strum(serialize = "s3:object-lock-retain-until-date")]
S3ObjectLockRetainUntilDate,
#[strum(serialize = "s3:object-lock-legal-hold")]
S3ObjectLockLegalHold,
#[strum(serialize = "s3:object-lock-mode")]
S3ObjectLockMode,
@@ -389,6 +392,7 @@ mod tests {
#[test_case("s3:VersionId", KeyName::S3(S3KeyName::S3VersionId) ; "aws_version_id")]
#[test_case("s3:versionid", KeyName::S3(S3KeyName::S3VersionId) ; "minio_version_id")]
#[test_case("s3:object-lock-mode", KeyName::S3(S3KeyName::S3ObjectLockMode))]
#[test_case("s3:object-lock-legal-hold", KeyName::S3(S3KeyName::S3ObjectLockLegalHold))]
#[test_case("aws:SecureTransport", KeyName::Aws(AwsKeyName::AWSSecureTransport))]
#[test_case("jwt:sub", KeyName::Jwt(JwtKeyName::JWTSub))]
#[test_case("ldap:user", KeyName::Ldap(LdapKeyName::User))]
@@ -412,6 +416,7 @@ mod tests {
#[test_case("s3:VersionId", KeyName::S3(S3KeyName::S3VersionId) ; "aws_version_id")]
#[test_case("s3:versionid", KeyName::S3(S3KeyName::S3VersionId) ; "minio_version_id")]
#[test_case("s3:object-lock-mode", KeyName::S3(S3KeyName::S3ObjectLockMode))]
#[test_case("s3:object-lock-legal-hold", KeyName::S3(S3KeyName::S3ObjectLockLegalHold))]
#[test_case("aws:SecureTransport", KeyName::Aws(AwsKeyName::AWSSecureTransport))]
#[test_case("jwt:sub", KeyName::Jwt(JwtKeyName::JWTSub))]
#[test_case("ldap:user", KeyName::Ldap(LdapKeyName::User))]
@@ -431,6 +436,7 @@ mod tests {
#[test_case("s3:x-amz-copy-source", KeyName::S3(S3KeyName::S3XAmzCopySource))]
#[test_case("s3:versionid", KeyName::S3(S3KeyName::S3VersionId))]
#[test_case("s3:object-lock-mode", KeyName::S3(S3KeyName::S3ObjectLockMode))]
#[test_case("s3:object-lock-legal-hold", KeyName::S3(S3KeyName::S3ObjectLockLegalHold))]
#[test_case("aws:SecureTransport", KeyName::Aws(AwsKeyName::AWSSecureTransport))]
#[test_case("jwt:sub", KeyName::Jwt(JwtKeyName::JWTSub))]
#[test_case("ldap:user", KeyName::Ldap(LdapKeyName::User))]
+2 -1
View File
@@ -287,7 +287,7 @@ mod tests {
};
use std::collections::HashMap;
use crate::policy::function::key_name::S3KeyName::{S3LocationConstraint, S3ObjectLockMode};
use crate::policy::function::key_name::S3KeyName::{S3LocationConstraint, S3ObjectLockLegalHold, S3ObjectLockMode};
use test_case::test_case;
fn new_func(name: KeyName, variable: Option<String>, values: Vec<&str>) -> StringFunc {
@@ -309,6 +309,7 @@ mod tests {
#[test_case(r#"{"aws:username/value": ["johndoe", "aaa"]}"#, new_func(Aws(AWSUsername), Some("value".into()), vec!["johndoe", "aaa"]
))]
#[test_case(r#"{"s3:object-lock-mode": "COMPLIANCE"}"#, new_func(S3(S3ObjectLockMode), None, vec!["COMPLIANCE"]))]
#[test_case(r#"{"s3:object-lock-legal-hold": "ON"}"#, new_func(S3(S3ObjectLockLegalHold), None, vec!["ON"]))]
fn test_deser(input: &str, expect: StringFunc) -> Result<(), serde_json::Error> {
let v: StringFunc = serde_json::from_str(input)?;
assert_eq!(v, expect);
+33
View File
@@ -131,6 +131,39 @@ pub(crate) async fn read_config_with_revision<S: ScannerObjectIO>(
}
}
pub(crate) fn usage_floor_primary_read_error_allows_backup(err: &Error) -> bool {
match err {
Error::FileCorrupt
| Error::CorruptedFormat
| Error::CorruptedBackend
| Error::PartMissingOrCorrupt
| Error::LessData
| Error::MoreData => true,
Error::Io(io_error) => {
matches!(io_error.kind(), std::io::ErrorKind::InvalidData | std::io::ErrorKind::UnexpectedEof)
|| error_chain_has_usage_floor_corruption_signature(io_error)
}
_ => false,
}
}
fn error_chain_has_usage_floor_corruption_signature(error: &(dyn std::error::Error + 'static)) -> bool {
let mut current = Some(error);
while let Some(err) = current {
let message = err.to_string();
if message.contains("InlineData value out of range")
|| message.contains("InlineData key out of range")
|| message.contains("insufficient data for metadata")
|| message.contains("insufficient data for meta length")
|| message.contains("insufficient data for CRC")
{
return true;
}
current = err.source();
}
false
}
/// Read only the object revision without materializing its body.
pub(crate) async fn read_config_revision<S: ScannerObjectIO>(store: Arc<S>, path: &str) -> StorageResult<DataUsageCacheRevision> {
match store
+18 -1
View File
@@ -14,7 +14,9 @@
/// Scanner cycle-state codec, persisted usage floors, and cycle-state persistence.
use super::*;
use crate::ScannerGetObjectReader;
use crate::data_usage_define::{DATA_USAGE_BLOOM_RECOVERY_PATH, DATA_USAGE_RECOVERY_PATH};
use crate::data_usage_define::{
DATA_USAGE_BLOOM_RECOVERY_PATH, DATA_USAGE_RECOVERY_PATH, usage_floor_primary_read_error_allows_backup,
};
use crate::storage_api::owner::ObjectIO as _;
use tokio::io::AsyncReadExt as _;
@@ -1717,6 +1719,7 @@ pub(super) async fn persisted_usage_floor_for_startup(
let backup_path = format!("{primary_path}.bkp");
let is_v2_path = primary_path == DATA_USAGE_OBJ_NAME_PATH.as_str();
let mut recovered_primary_companion_epoch = None;
let mut primary_read_error = None;
let primary_epoch = match read_config_with_revision(storeapi.clone(), primary_path).await {
Ok((Some(data), revision)) => {
let usage = serde_json::from_slice::<DataUsageInfo>(&data).map_err(|err| {
@@ -1776,6 +1779,12 @@ pub(super) async fn persisted_usage_floor_for_startup(
}
}
Ok((None, _)) => None,
Err(err) if !is_v2_path && usage_floor_primary_read_error_allows_backup(&err) => {
primary_read_error = Some(format!("failed to read scanner usage epoch floor from {primary_path}: {err}"));
invalid_baseline_path.get_or_insert_with(|| primary_path.to_string());
unrecoverable_baseline_path.get_or_insert_with(|| primary_path.to_string());
None
}
Err(err) => {
return Err(ScannerError::Other(format!(
"failed to read scanner usage epoch floor from {primary_path}: {err}"
@@ -1855,6 +1864,14 @@ pub(super) async fn persisted_usage_floor_for_startup(
)));
}
}
if let Some(primary_read_error) = primary_read_error
&& !any_found
{
return Err(ScannerError::Other(format!(
"{}; no valid scanner usage floor backup was available at {backup_path}",
primary_read_error
)));
}
if any_found {
if bootstrap_pending {
return Err(ScannerError::Other(
+29 -8
View File
@@ -13,6 +13,7 @@
// limitations under the License.
/// Leader-lock claiming, usage-epoch fencing, and lock-loss handling.
use super::*;
use crate::data_usage_define::usage_floor_primary_read_error_allows_backup;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum ScannerLeadershipClaimReconcile {
@@ -142,20 +143,34 @@ pub(super) async fn usage_snapshot_for_epoch_fence(
}
}
for path in [
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
] {
let (legacy, _) = read_config_with_revision(storeapi.clone(), &path)
.await
.map_err(|err| ScannerError::Other(format!("failed to read legacy scanner usage epoch fence: {err}")))?;
let legacy_primary_path = LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string();
let legacy_backup_path = format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str());
let mut legacy_primary_read_error = None;
for path in [&legacy_primary_path, &legacy_backup_path] {
let legacy = match read_config_with_revision(storeapi.clone(), path).await {
Ok((legacy, _)) => legacy,
Err(err) if path == &legacy_primary_path && usage_floor_primary_read_error_allows_backup(&err) => {
legacy_primary_read_error = Some(format!("failed to read legacy scanner usage epoch fence from {path}: {err}"));
continue;
}
Err(err) => {
return Err(ScannerError::Other(format!(
"failed to read legacy scanner usage epoch fence from {path}: {err}"
)));
}
};
if let Some(legacy) = legacy.as_deref() {
let usage = decode_usage_snapshot_for_epoch_fence(legacy, &path, false)?;
let usage = decode_usage_snapshot_for_epoch_fence(legacy, path, false)?;
if invalid_primary_epoch.is_none_or(|epoch| usage.scanner_epoch.unwrap_or_default() >= epoch) {
return Ok(Some(usage));
}
}
}
if let Some(legacy_primary_read_error) = legacy_primary_read_error {
return Err(ScannerError::Other(format!(
"{legacy_primary_read_error}; no valid legacy scanner usage epoch fence backup was available at {legacy_backup_path}"
)));
}
// A missing usage snapshot is an uninitialized state, not an empty
// snapshot. Leadership fencing may proceed without creating a plausible
// default; the first authoritative scanner publication will create it.
@@ -258,6 +273,12 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
Some(epoch) if epoch == claimed_epoch => return Ok(()),
Some(_) | None => {}
}
// A validated pre-marker legacy baseline needs an explicit complete
// identity before acquiring an epoch. Otherwise the v2 reader would
// reject the fenced value on its next startup.
if !usage.usage_snapshot_bootstrap_pending {
usage.usage_snapshot_complete = true;
}
usage.scanner_epoch = Some(claimed_epoch);
let data = serde_json::to_vec(&usage)
.map_err(|err| ScannerError::Other(format!("failed to encode scanner usage epoch fence: {err}")))?;
+314
View File
@@ -574,6 +574,7 @@ struct MemoryConfigStore {
objects: Mutex<HashMap<String, Vec<u8>>>,
revisions: Mutex<HashMap<String, u64>>,
insert_after_gets: Mutex<HashMap<String, Vec<u8>>>,
read_errors: Mutex<HashMap<String, EcstoreError>>,
delayed_gets: Mutex<HashMap<String, Duration>>,
non_regular_objects: Mutex<HashSet<String>>,
fail_put_number: Mutex<HashMap<String, usize>>,
@@ -623,6 +624,9 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore {
_opts: &ObjectOptions,
) -> EcstoreResult<GetObjectReader> {
let key = memory_config_key(bucket, object);
if let Some(error) = self.read_errors.lock().await.get(&key).cloned() {
return Err(error);
}
if let Some(delay) = self.delayed_gets.lock().await.remove(&key) {
tokio::time::sleep(delay).await;
}
@@ -2902,6 +2906,147 @@ async fn scanner_usage_floor_recovers_from_incomplete_v2_primary_using_fenced_ba
);
}
async fn seed_legacy_primary_read_error_with_backup(store: &Arc<MemoryConfigStore>, error: EcstoreError, epoch: u64, cycle: u64) {
let legacy_primary = LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str();
let legacy_backup = format!("{legacy_primary}.bkp");
let mut backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
backup.scanner_epoch = Some(epoch);
backup.scanner_cycle = Some(cycle);
store
.read_errors
.lock()
.await
.insert(memory_config_key(RUSTFS_META_BUCKET, legacy_primary), error);
store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, &legacy_backup),
serde_json::to_vec(&backup).expect("legacy backup usage snapshot should encode"),
);
}
#[tokio::test]
async fn scanner_usage_floor_recovers_legacy_backup_after_primary_decode_error() {
let store = Arc::new(MemoryConfigStore::default());
seed_legacy_primary_read_error_with_backup(&store, EcstoreError::other("InlineData value out of range"), 19, 41).await;
let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true)
.await
.expect("valid legacy backup should recover the startup floor");
assert_eq!(state, PersistedUsageFloorStartup::Authoritative);
assert_eq!(
floor,
PersistedUsageFloor {
next_cycle: 42,
leader_epoch: 19,
}
);
assert_eq!(
persisted_usage_floor(store)
.await
.expect("valid legacy backup should recover the authoritative floor"),
floor
);
}
#[tokio::test]
async fn scanner_usage_floor_does_not_bootstrap_over_corrupt_legacy_primary_without_backup() {
let store = Arc::new(MemoryConfigStore::default());
let legacy_primary = LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str();
store
.read_errors
.lock()
.await
.insert(memory_config_key(RUSTFS_META_BUCKET, legacy_primary), EcstoreError::FileCorrupt);
let err = persisted_usage_floor_for_startup(store, true)
.await
.expect_err("corrupt legacy primary without a valid backup must remain fail-closed");
assert!(err.to_string().contains("no valid scanner usage floor backup"), "unexpected error: {err}");
}
#[tokio::test]
async fn scanner_usage_floor_does_not_fallback_to_legacy_after_corrupt_v2_primary() {
let store = Arc::new(MemoryConfigStore::default());
let mut v2_backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
v2_backup.scanner_epoch = Some(8);
v2_backup.scanner_cycle = Some(11);
let mut legacy = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
legacy.scanner_epoch = Some(3);
legacy.scanner_cycle = Some(7);
store.read_errors.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
EcstoreError::FileCorrupt,
);
store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, &format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str())),
serde_json::to_vec(&v2_backup).expect("v2 backup usage snapshot should encode"),
);
store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
serde_json::to_vec(&legacy).expect("legacy usage snapshot should encode"),
);
let err = persisted_usage_floor_for_startup(store, true)
.await
.expect_err("corrupt v2 primary must not recover without a primary revision");
assert!(
err.to_string().contains(&format!(
"failed to read scanner usage epoch floor from {}",
DATA_USAGE_OBJ_NAME_PATH.as_str()
)),
"unexpected error: {err}"
);
}
#[tokio::test]
async fn scanner_usage_floor_keeps_transient_primary_read_error_fail_closed() {
let store = Arc::new(MemoryConfigStore::default());
let legacy_primary = LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str();
seed_legacy_primary_read_error_with_backup(
&store,
EcstoreError::Io(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"connection reset while reading usage primary",
)),
19,
41,
)
.await;
let err = persisted_usage_floor_for_startup(store, true)
.await
.expect_err("transient primary errors must not be converted into backup recovery");
assert!(
err.to_string()
.contains(&format!("failed to read scanner usage epoch floor from {legacy_primary}")),
"unexpected error: {err}"
);
assert!(
!err.to_string().contains("no valid scanner usage floor backup"),
"transient error should not enter corrupt-primary fallback: {err}"
);
}
#[tokio::test]
async fn scanner_usage_floor_keeps_outdated_primary_metadata_fail_closed() {
let store = Arc::new(MemoryConfigStore::default());
let legacy_primary = LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str();
seed_legacy_primary_read_error_with_backup(&store, EcstoreError::OutdatedXLMeta, 19, 41).await;
let err = persisted_usage_floor_for_startup(store, true)
.await
.expect_err("outdated primary metadata must not be converted into backup recovery");
assert!(
err.to_string()
.contains(&format!("failed to read scanner usage epoch floor from {legacy_primary}")),
"unexpected error: {err}"
);
assert!(
!err.to_string().contains("no valid scanner usage floor backup"),
"outdated metadata should not enter corrupt-primary fallback: {err}"
);
}
#[tokio::test]
async fn scanner_usage_floor_does_not_bootstrap_over_incomplete_v2_primary() {
let store = Arc::new(MemoryConfigStore::default());
@@ -3004,6 +3149,19 @@ async fn scanner_leadership_fencing_recovers_incomplete_v2_primary_from_backup()
assert_eq!(recovered.scanner_cycle, Some(103));
}
#[tokio::test]
async fn scanner_usage_floor_leadership_fencing_recovers_legacy_backup_after_primary_decode_error() {
let store = Arc::new(MemoryConfigStore::default());
seed_legacy_primary_read_error_with_backup(&store, EcstoreError::other("InlineData value out of range"), 19, 41).await;
let recovered = usage_snapshot_for_epoch_fence(store, None, false)
.await
.expect("a valid legacy backup should provide the fencing baseline")
.expect("the fencing baseline should be present");
assert_eq!(recovered.scanner_epoch, Some(19));
assert_eq!(recovered.scanner_cycle, Some(41));
}
#[tokio::test]
async fn scanner_usage_floor_ignores_older_backup_after_primary_epoch_fence() {
let store = Arc::new(MemoryConfigStore::default());
@@ -3955,6 +4113,162 @@ async fn scanner_defers_leadership_when_usage_snapshots_are_stably_absent() {
assert!(read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()).await.is_err());
}
#[tokio::test]
async fn scanner_usage_floor_leadership_claim_recovers_legacy_backup_after_primary_decode_error() {
let store = Arc::new(MemoryConfigStore::default());
let ctx = CancellationToken::new();
seed_legacy_primary_read_error_with_backup(&store, EcstoreError::other("InlineData value out of range"), 19, 41).await;
let mut revision = DataUsageCacheRevision::Missing;
let mut cycle = CurrentCycle::default();
let mut persisted_epoch = 19;
assert!(
claim_scanner_leadership(
&ctx,
store.clone(),
&mut cycle,
&mut revision,
&mut persisted_epoch,
false,
ScannerCycleResetPolicy::None,
)
.await
);
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("leadership claim should persist after legacy backup recovery");
let (_, claimed_epoch) = decode_scanner_cycle_state(&state).expect("leadership claim should decode");
assert_eq!(claimed_epoch, 20);
assert_eq!(persisted_epoch, 20);
let usage = read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("legacy backup recovery should publish a fenced v2 usage primary");
let usage = serde_json::from_slice::<DataUsageInfo>(&usage).expect("fenced v2 usage primary should decode");
assert_eq!(usage.scanner_epoch, Some(20));
assert_eq!(usage.scanner_cycle, Some(41));
}
#[tokio::test]
#[serial_test::serial]
async fn scanner_legacy_usage_backup_survives_fencing_and_restart_after_real_metadata_truncation() {
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
let (temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
usage.usage_snapshot_complete = false;
usage.scanner_cycle = Some(41);
let mut data = serde_json::to_vec(&usage).expect("legacy usage should encode");
data.resize(data.len() + 16 * 1024, b' ');
let legacy_path = LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str();
let backup_path = format!("{legacy_path}.bkp");
for path in [legacy_path, backup_path.as_str()] {
save_config(store.clone(), path, data.clone())
.await
.expect("legacy usage fixture should persist");
}
let mut truncated_files = Vec::new();
for disk_index in 0..4 {
let path = temp_dir
.path()
.join(format!("pool0/disk{disk_index}"))
.join(RUSTFS_META_BUCKET)
.join(legacy_path)
.join("xl.meta");
let file = tokio::fs::OpenOptions::new()
.write(true)
.open(&path)
.await
.expect("legacy inline metadata should exist");
assert!(file.metadata().await.expect("metadata should be readable").len() > 4096);
file.set_len(4096).await.expect("fixture should truncate at a page boundary");
truncated_files.push((
path.clone(),
tokio::fs::read(&path)
.await
.expect("truncated evidence should remain readable"),
));
}
let store = restart_scanner_cycle_store_from(&store).await;
let error = read_config_with_revision(store.clone(), legacy_path)
.await
.expect_err("truncated primary must fail in the real object reader");
assert!(
error.to_string().contains("InlineData value out of range"),
"unexpected truncated-primary error: {error}"
);
assert_eq!(
read_config_with_revision(store.clone(), &backup_path)
.await
.expect("backup should remain readable")
.0,
Some(data.clone()),
);
let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true)
.await
.expect("intact legacy backup must recover startup despite truncated primary");
assert_eq!(
floor,
PersistedUsageFloor {
next_cycle: 42,
leader_epoch: 0
}
);
assert_eq!(state, PersistedUsageFloorStartup::Authoritative);
let baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("publication must also read the intact backup");
assert_eq!(baseline.data.as_deref(), Some(data.as_slice()));
assert_eq!(baseline.revision, DataUsageCacheRevision::Missing);
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false)
.await
.expect("legacy backup must be fenced into v2");
let fenced = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("fencing must publish a v2 usage primary");
let fenced = serde_json::from_slice::<DataUsageInfo>(&fenced).expect("fenced v2 usage primary should decode");
assert!(
fenced.usage_snapshot_complete,
"the fenced pre-marker baseline must become a complete v2 identity"
);
let store = restart_scanner_cycle_store_from(&store).await;
let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true)
.await
.expect("a restart after fencing must preserve the recovered floor");
assert_eq!(
floor,
PersistedUsageFloor {
next_cycle: 42,
leader_epoch: 7
}
);
assert_eq!(state, PersistedUsageFloorStartup::Authoritative);
let restarted = restart_scanner_cycle_store_from(&store).await;
assert_eq!(
persisted_usage_floor(restarted)
.await
.expect("fenced floor must survive another restart"),
PersistedUsageFloor {
next_cycle: 42,
leader_epoch: 7
}
);
for (path, bytes) in truncated_files {
assert_eq!(tokio::fs::read(path).await.expect("legacy evidence must not be removed"), bytes);
}
assert_eq!(
read_config(store, &backup_path)
.await
.expect("legacy backup must remain intact"),
data
);
global_metrics().set_cycle(None).await;
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
async fn usage_bootstrap_pending_unblocks_first_leadership_claim() {
let store = Arc::new(MemoryConfigStore::default());
+27 -2
View File
@@ -13,6 +13,7 @@
// limitations under the License.
/// Data-usage snapshot persistence: CAS store pipeline, epoch baselines, and observed-snapshot cleanup.
use super::*;
use crate::data_usage_define::usage_floor_primary_read_error_allows_backup;
use crate::storage_api::owner::ScannerPublicationCommitState;
use std::collections::HashMap;
use std::sync::atomic::AtomicBool;
@@ -53,6 +54,30 @@ pub(super) struct DataUsagePersistBaseline {
pub(super) revision: DataUsageCacheRevision,
}
async fn read_usage_persist_candidate(
storeapi: Arc<impl ScannerObjectIO>,
path: &str,
) -> Result<(Option<Vec<u8>>, DataUsageCacheRevision), EcstoreError> {
let primary = read_config_with_revision(storeapi.clone(), path).await;
if path != LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str() {
return primary;
}
let Err(primary_error) = &primary else {
return primary;
};
if !usage_floor_primary_read_error_allows_backup(primary_error) {
return primary;
}
let backup_path = format!("{path}.bkp");
let backup = read_config_with_revision(storeapi, &backup_path).await?;
if backup.0.as_deref().is_some_and(|data| {
serde_json::from_slice::<DataUsageInfo>(data).is_ok_and(|usage| data_usage_info_has_persisted_baseline_identity(&usage))
}) {
return Ok(backup);
}
primary
}
/// Read the bytes used as the baseline for a usage publication while keeping
/// the v2 primary revision as the CAS fence. During an interrupted upgrade the
/// primary can be valid JSON without a baseline identity; in that case a
@@ -68,7 +93,7 @@ pub(super) async fn read_data_usage_persist_baseline(
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
] {
let (candidate, _) = read_config_with_revision(storeapi.clone(), &path).await?;
let (candidate, _) = read_usage_persist_candidate(storeapi.clone(), &path).await?;
let Some(candidate) = candidate else {
continue;
};
@@ -107,7 +132,7 @@ pub(super) async fn read_data_usage_persist_baseline(
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
] {
let (candidate, _) = read_config_with_revision(storeapi.clone(), &path).await?;
let (candidate, _) = read_usage_persist_candidate(storeapi.clone(), &path).await?;
let Some(candidate) = candidate else {
continue;
};
+1 -1
View File
@@ -77,7 +77,7 @@
rustfs = rustPlatform.buildRustPackage {
pname = "rustfs";
version = "1.0.0-rc.4";
version = "1.0.0-rc.5";
src = ./.;
+2 -2
View File
@@ -2,8 +2,8 @@ apiVersion: v2
name: rustfs
description: RustFS helm chart to deploy RustFS on kubernetes cluster.
type: application
version: "1.0.0-rc.4"
appVersion: "1.0.0-rc.4"
version: "1.0.0-rc.5"
appVersion: "1.0.0-rc.5"
home: https://rustfs.com
icon: https://media.sys.truenas.net/apps/rustfs/icons/icon.svg
maintainers:
+5 -2
View File
@@ -1,9 +1,9 @@
%global _enable_debug_packages 0
%global _empty_manifest_terminate_build 0
%global prerelease rc.4
%global prerelease rc.5
Name: rustfs
Version: 1.0.0
Release: rc.4
Release: rc.5
Summary: High-performance distributed object storage for MinIO alternative
License: Apache-2.0
@@ -58,6 +58,9 @@ install %_builddir/%{name}-%{version}-%{prerelease}/target/%_arch/%_arch-unknown
%_bindir/rustfs
%changelog
* Mon Aug 31 2026 overtrue <[email protected]>
- Update RPM package to RustFS 1.0.0-rc.5
* Thu Aug 27 2026 overtrue <[email protected]>
- Update RPM package to RustFS 1.0.0-rc.4
+401 -56
View File
@@ -32,7 +32,9 @@ use crate::admin::storage_api::bucket::replication::{
};
use crate::admin::storage_api::bucket::target::{BucketTarget, BucketTargetType, BucketTargets};
use crate::admin::storage_api::bucket::utils::{deserialize, serialize};
use crate::admin::storage_api::bucket::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _};
use crate::admin::storage_api::bucket::{
AdminObjectLockConfigExt as _, AdminReplicationConfigExt as _, AdminVersioningConfigExt as _,
};
use crate::admin::storage_api::contract::bucket::{
BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp,
};
@@ -61,18 +63,18 @@ use rustfs_madmin::{
BucketBandwidth, GroupStatus, IDPSettings, InProgressMetric, InQueueMetric, LDAPConfigSettings, LDAPSettings,
OpenIDProviderSettings, PeerInfo, PeerSite, QStat, ReplProxyMetric, ReplicateAddStatus, ReplicateEditStatus,
ReplicateRemoveStatus, ResyncBucketStatus, SITE_REPL_API_VERSION, SR_IAM_ITEM_STS_ACC, SR_IAM_ITEM_STS_ACC_LEGACY,
SRBucketMeta, SRBucketStatsSummary, SRGroupInfo, SRGroupStatsSummary, SRIAMItem, SRIAMUser, SRILMExpiryStatsSummary, SRInfo,
SRMetric, SRMetricsSummary, SRPeerError, SRPeerJoinReq, SRPendingOperation, SRPolicyMapping, SRPolicyStatsSummary,
SRRemoveReq, SRResyncOpStatus, SRSTSCredential, SRSessionPolicy, SRSiteSummary, SRStateEditReq, SRStateInfo, SRStatusInfo,
SRSvcAccChange, SRSvcAccCreate, SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
SRBucketInfo, SRBucketMeta, SRBucketStatsSummary, SRGroupInfo, SRGroupStatsSummary, SRIAMItem, SRIAMUser,
SRILMExpiryStatsSummary, SRInfo, SRMetric, SRMetricsSummary, SRPeerError, SRPeerJoinReq, SRPendingOperation, SRPolicyMapping,
SRPolicyStatsSummary, SRRemoveReq, SRResyncOpStatus, SRSTSCredential, SRSessionPolicy, SRSiteSummary, SRStateEditReq,
SRStateInfo, SRStatusInfo, SRSvcAccChange, SRSvcAccCreate, SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
};
use rustfs_policy::policy::{
Policy,
action::{Action, AdminAction},
};
use s3s::dto::{
DeleteMarkerReplicationStatus, DeleteReplicationStatus, ExistingObjectReplicationStatus, ReplicaModificationsStatus,
ReplicationConfiguration, ReplicationRule, ReplicationRuleStatus,
DeleteMarkerReplicationStatus, DeleteReplicationStatus, ExistingObjectReplicationStatus, ObjectLockConfiguration,
ReplicaModificationsStatus, ReplicationConfiguration, ReplicationRule, ReplicationRuleStatus, VersioningConfiguration,
};
use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
use serde::Deserialize;
@@ -287,12 +289,21 @@ struct SiteReplicationAddPreflightInfo {
endpoint: String,
deployment_id: String,
enabled: bool,
bucket_count: usize,
bucket_names: HashSet<String>,
buckets: BTreeMap<String, AddPreflightBucketCompat>,
peer_deployment_ids: BTreeSet<String>,
idp_settings: serde_json::Value,
}
/// The per-bucket facts the add preflight compares across sites when more
/// than one requested site holds data (rustfs/backlog#2070). Only properties
/// that cannot converge after the add belong here — everything else is
/// reconciled by the bucket-metadata sync.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct AddPreflightBucketCompat {
versioning_enabled: bool,
object_lock_enabled: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
struct SRPeerJoinResponse {
peer: PeerInfo,
@@ -935,19 +946,65 @@ fn idp_settings_value(settings: &IDPSettings) -> S3Result<serde_json::Value> {
.map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize IDP settings failed: {e}")))
}
/// The merge-critical facts of one bucket a site reported in its add
/// preflight metainfo. The configs arrive as the `build_sr_info` wire form
/// (base64-encoded XML); an undecodable config fails the preflight instead of
/// defaulting, so corruption cannot admit an unsafe merge.
fn add_preflight_bucket_compat(endpoint: &str, bucket: &str, info: &SRBucketInfo) -> S3Result<AddPreflightBucketCompat> {
let versioning_enabled = info
.versioning
.as_deref()
.map(|raw| {
deserialize::<VersioningConfiguration>(&decode_bucket_meta_wire_value(raw))
.map(|config| config.enabled())
.map_err(|e| {
s3_error!(
InvalidRequest,
"site `{endpoint}` reported an unreadable versioning config for bucket `{bucket}`: {e}"
)
})
})
.transpose()?
.unwrap_or(false);
let object_lock_enabled = info
.object_lock_config
.as_deref()
.map(|raw| {
deserialize::<ObjectLockConfiguration>(&decode_bucket_meta_wire_value(raw))
.map(|config| config.enabled())
.map_err(|e| {
s3_error!(
InvalidRequest,
"site `{endpoint}` reported an unreadable object-lock config for bucket `{bucket}`: {e}"
)
})
})
.transpose()?
.unwrap_or(false);
Ok(AddPreflightBucketCompat {
versioning_enabled,
object_lock_enabled,
})
}
fn add_preflight_info_from_sr_info(
site: &PeerSite,
info: SRInfo,
idp_settings: IDPSettings,
) -> S3Result<SiteReplicationAddPreflightInfo> {
let bucket_names = info.buckets.keys().cloned().collect();
let buckets = info
.buckets
.iter()
.map(|(bucket, bucket_info)| {
add_preflight_bucket_compat(&site.endpoint, bucket, bucket_info).map(|compat| (bucket.clone(), compat))
})
.collect::<S3Result<BTreeMap<_, _>>>()?;
Ok(SiteReplicationAddPreflightInfo {
name: if info.name.is_empty() { site.name.clone() } else { info.name },
endpoint: site.endpoint.clone(),
deployment_id: info.deployment_id,
enabled: info.enabled,
bucket_count: info.buckets.len(),
bucket_names,
buckets,
peer_deployment_ids: info.state.peers.keys().cloned().collect(),
idp_settings: idp_settings_value(&idp_settings)?,
})
@@ -1045,7 +1102,7 @@ fn validate_add_preflight_topology(infos: &[SiteReplicationAddPreflightInfo], lo
if info.deployment_id == local_peer.deployment_id {
local_seen = true;
}
if info.bucket_count > 0 {
if !info.buckets.is_empty() {
non_empty_sites.push(info.name.clone());
}
}
@@ -1070,11 +1127,15 @@ fn validate_add_preflight_topology(infos: &[SiteReplicationAddPreflightInfo], lo
}
if non_empty_sites.len() > 1 {
return Err(s3_error!(
InvalidRequest,
"site replication can be initialized with data on only one site; non-empty sites: {}",
non_empty_sites.join(", ")
));
validate_nonempty_add_bucket_compatibility(infos)?;
info!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "nonempty_sites_admitted",
non_empty_sites = %non_empty_sites.join(", "),
"admin site replication state"
);
}
let requested: BTreeSet<String> = infos.iter().map(|info| info.deployment_id.clone()).collect();
@@ -1091,6 +1152,76 @@ fn validate_add_preflight_topology(infos: &[SiteReplicationAddPreflightInfo], lo
Ok(())
}
/// Operator recovery guidance for a rejected add between sites that both hold
/// data — the only supported path is to empty one side and let a resync copy
/// the objects back (rustfs/backlog#2070).
const NONEMPTY_ADD_RECOVERY_HINT: &str = "to pair these sites, delete the conflicting bucket (or its data) on all but one \
site, re-run `replicate add`, then run `replicate resync` from the surviving site to restore the objects";
/// Admission check for an add in which more than one requested site holds
/// data — the DR re-pair case: two sites that were unpaired (or never
/// finished a removal) both keep their buckets, and the historical
/// unconditional "only one site may hold data" rejection made `replicate
/// remove` a one-way door (rustfs/backlog#2070).
///
/// The add is admitted when every bucket name held by MORE than one requested
/// site is provably safe to merge through the existing backfill/resync
/// convergence:
///
/// - versioning must be Enabled on every holder: replication into a versioned
/// bucket lands as another version, so a same-key object from the peer
/// never destroys the local copy — while on an unversioned holder it would
/// silently replace the only copy;
/// - object-lock enablement must match across holders: lock cannot be toggled
/// after bucket creation, so a mismatch never converges, and replicating
/// locked objects into a lock-less bucket would strip their WORM guarantee.
///
/// A bucket held by a single site carries no merge risk — the post-add
/// backfill creates it on the peers exactly as the historical
/// one-non-empty-site path always has.
fn validate_nonempty_add_bucket_compatibility(infos: &[SiteReplicationAddPreflightInfo]) -> S3Result<()> {
let mut holders: BTreeMap<&str, Vec<(&SiteReplicationAddPreflightInfo, AddPreflightBucketCompat)>> = BTreeMap::new();
for info in infos {
for (bucket, compat) in &info.buckets {
holders.entry(bucket.as_str()).or_default().push((info, *compat));
}
}
for (bucket, holders) in holders {
let [(first, first_compat), rest @ ..] = holders.as_slice() else {
continue;
};
if rest.is_empty() {
continue;
}
if let Some((conflicting, _)) = rest
.iter()
.find(|(_, compat)| compat.object_lock_enabled != first_compat.object_lock_enabled)
{
let (enabled_on, disabled_on) = if first_compat.object_lock_enabled {
(&first.name, &conflicting.name)
} else {
(&conflicting.name, &first.name)
};
return Err(s3_error!(
InvalidRequest,
"bucket `{bucket}` has object lock enabled on site `{enabled_on}` but not on site `{disabled_on}`, and \
object lock cannot be changed after bucket creation; {NONEMPTY_ADD_RECOVERY_HINT}"
));
}
if let Some((unversioned, _)) = holders.iter().find(|(_, compat)| !compat.versioning_enabled) {
return Err(s3_error!(
InvalidRequest,
"bucket `{bucket}` exists on more than one site but does not have versioning enabled on site `{}`, so \
merging could silently overwrite objects; enable versioning on every site holding it, or {NONEMPTY_ADD_RECOVERY_HINT}",
unversioned.name
));
}
}
Ok(())
}
fn site_replication_bootstrap_token(uri: &Uri) -> Option<String> {
query_pairs(uri).get("bootstrapToken").cloned()
}
@@ -5701,6 +5832,53 @@ fn sts_replication_compatibility_policy<'a>(claims: &HashMap<String, Value>, par
(!claims.contains_key(OIDC_VIRTUAL_PARENT_CLAIM) && !parent_policy_mapping.is_empty()).then_some(parent_policy_mapping)
}
/// Adopt only the fields a committed add computed onto the freshly loaded
/// transaction state. Everything else is owned by writers that commit without
/// touching `updated_at` (retry events, peer-edit generations, resync
/// progress, the acks/clears of an already pending rotation), so the add's
/// `updated_at` CAS cannot vouch for them — they keep the freshly loaded
/// value, except `pending_remove`:
///
/// A committed add supersedes a half-finished removal THIS site started,
/// exactly as an accepted join does on the receiving side (`apply_peer_join`,
/// rustfs/rustfs#5963): the adopted topology IS the new membership, while the
/// pending record only exists to keep notifying peers about the old one. Left
/// in place, the reconcile tick would replay the stale `SRRemoveReq` against
/// a freshly re-paired peer — `SRPeerRemoveHandler` applies it
/// unconditionally — and dismantle the pairing this add just created
/// (rustfs/backlog#2070). A removal that started AFTER the add's preflight
/// snapshot moved `updated_at`, so the CAS refuses the commit before this
/// runs.
///
/// The exhaustive destructure makes adding a state field a compile error here
/// until it is classified.
fn adopt_add_commit_state(state: &mut SiteReplicationState, next_state: SiteReplicationState) {
let SiteReplicationState {
name,
service_account_access_key,
service_account_secret_key: _,
service_account_parent,
peers,
updated_at,
resync_status: _,
pending_rotation: _,
pending_remove: _,
pending_endpoint_refresh: _,
retry_queue: _,
iam_deletion_replays: _,
sync_state_initialized,
edit_generation: _,
applied_edit_generations: _,
} = next_state;
state.name = name;
state.service_account_access_key = service_account_access_key;
state.service_account_parent = service_account_parent;
state.peers = peers;
state.updated_at = updated_at;
state.sync_state_initialized = sync_state_initialized;
state.pending_remove = None;
}
pub struct SiteReplicationAddHandler {}
/// MinIO's `SRPeerJoin` replies with an empty body on success; synthesize the
@@ -5756,7 +5934,7 @@ impl Operation for SiteReplicationAddHandler {
let bootstrap_buckets = preflight_infos
.iter()
.filter(|info| !same_identity_endpoint(&info.endpoint, &local_peer.endpoint))
.flat_map(|info| info.bucket_names.iter().cloned())
.flat_map(|info| info.buckets.keys().cloned())
.collect();
let add_in_progress_guard = SiteReplicationAddInProgressGuard::start(lifecycle_guard, bootstrap_buckets)?;
let mut state = merge_add_sites(
@@ -5853,36 +6031,7 @@ impl Operation for SiteReplicationAddHandler {
"site replication state changed during peer join; the peers may already be joined — re-run replicate add"
));
}
// Adopt only the fields this add computed. Everything else is
// owned by writers that commit without touching `updated_at`
// (retry events, peer-edit generations, resync progress, the
// acks/clears of an already pending rotation or removal), so the
// CAS above cannot vouch for them — they keep the freshly loaded
// value. The exhaustive destructure makes adding a state field a
// compile error here until it is classified.
let SiteReplicationState {
name,
service_account_access_key,
service_account_secret_key: _,
service_account_parent,
peers,
updated_at,
resync_status: _,
pending_rotation: _,
pending_remove: _,
pending_endpoint_refresh: _,
retry_queue: _,
iam_deletion_replays: _,
sync_state_initialized,
edit_generation: _,
applied_edit_generations: _,
} = next_state;
state.name = name;
state.service_account_access_key = service_account_access_key;
state.service_account_parent = service_account_parent;
state.peers = peers;
state.updated_at = updated_at;
state.sync_state_initialized = sync_state_initialized;
adopt_add_commit_state(state, next_state);
let edit_generation = next_peer_edit_generation(state);
Ok((state.clone(), edit_generation))
})
@@ -9421,18 +9570,29 @@ mod tests {
}
fn preflight_site(name: &str, endpoint: &str, deployment_id: &str, bucket_count: usize) -> SiteReplicationAddPreflightInfo {
// Site-prefixed names keep the generated buckets disjoint across
// sites; tests exercising shared-bucket merges insert their own.
let buckets = (0..bucket_count)
.map(|i| (format!("{name}-bucket-{i}"), versioned_bucket()))
.collect();
SiteReplicationAddPreflightInfo {
name: name.to_string(),
endpoint: endpoint.to_string(),
deployment_id: deployment_id.to_string(),
enabled: false,
bucket_count,
bucket_names: HashSet::new(),
buckets,
peer_deployment_ids: BTreeSet::new(),
idp_settings: serde_json::json!({"provider": "same"}),
}
}
fn versioned_bucket() -> AddPreflightBucketCompat {
AddPreflightBucketCompat {
versioning_enabled: true,
object_lock_enabled: false,
}
}
#[test]
fn test_validate_add_preflight_topology_accepts_matching_sites() {
let local_peer = PeerInfo {
@@ -9491,20 +9651,120 @@ mod tests {
assert!(err.to_string().contains("IDP settings mismatch"));
}
// rustfs/backlog#2070: two sites that both hold data (the DR re-pair
// case) must be admitted when their bucket sets are merge-safe, instead
// of the historical unconditional "only one site may hold data" rejection
// that made `replicate remove` a one-way door.
#[test]
fn test_validate_add_preflight_topology_rejects_multiple_non_empty_sites() {
fn test_validate_add_preflight_topology_accepts_compatible_non_empty_sites() {
let local_peer = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let mut local = preflight_site("local", "https://local.example.com", "local-dep", 1);
let mut remote = preflight_site("remote", "https://remote.example.com", "remote-dep", 1);
// The same bucket on both sites, versioning enabled on both: the
// exact shape a formerly paired cluster is left in after a remove.
local.buckets.insert("shared".to_string(), versioned_bucket());
remote.buckets.insert("shared".to_string(), versioned_bucket());
let infos = vec![local, remote];
validate_add_preflight_topology(&infos, &local_peer).expect("compatible non-empty sites should be admitted");
}
#[test]
fn test_validate_add_preflight_topology_accepts_disjoint_non_empty_sites() {
let local_peer = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let infos = vec![
preflight_site("local", "https://local.example.com", "local-dep", 1),
preflight_site("remote", "https://remote.example.com", "remote-dep", 1),
preflight_site("local", "https://local.example.com", "local-dep", 2),
preflight_site("remote", "https://remote.example.com", "remote-dep", 2),
];
let err = validate_add_preflight_topology(&infos, &local_peer).expect_err("multiple non-empty sites should fail");
validate_add_preflight_topology(&infos, &local_peer).expect("disjoint non-empty sites should be admitted");
}
assert!(err.to_string().contains("only one site"));
#[test]
fn test_validate_add_preflight_topology_rejects_shared_bucket_object_lock_mismatch() {
let local_peer = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let mut local = preflight_site("local", "https://local.example.com", "local-dep", 0);
let mut remote = preflight_site("remote", "https://remote.example.com", "remote-dep", 0);
local.buckets.insert(
"shared".to_string(),
AddPreflightBucketCompat {
versioning_enabled: true,
object_lock_enabled: true,
},
);
remote.buckets.insert("shared".to_string(), versioned_bucket());
let infos = vec![local, remote];
let err = validate_add_preflight_topology(&infos, &local_peer).expect_err("object-lock mismatch should fail");
let message = err.to_string();
assert!(
message.contains("bucket `shared` has object lock enabled on site `local`"),
"got: {message}"
);
// The rejection must carry the operator recovery steps, not a bare no.
assert!(message.contains("re-run `replicate add`"), "got: {message}");
assert!(message.contains("`replicate resync`"), "got: {message}");
}
#[test]
fn test_validate_add_preflight_topology_rejects_shared_unversioned_bucket() {
let local_peer = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let mut local = preflight_site("local", "https://local.example.com", "local-dep", 0);
let mut remote = preflight_site("remote", "https://remote.example.com", "remote-dep", 0);
local.buckets.insert("shared".to_string(), versioned_bucket());
remote.buckets.insert(
"shared".to_string(),
AddPreflightBucketCompat {
versioning_enabled: false,
object_lock_enabled: false,
},
);
let infos = vec![local, remote];
let err = validate_add_preflight_topology(&infos, &local_peer).expect_err("shared unversioned bucket should fail");
let message = err.to_string();
assert!(
message.contains("bucket `shared`") && message.contains("versioning enabled on site `remote`"),
"got: {message}"
);
assert!(message.contains("re-run `replicate add`"), "got: {message}");
}
// A bucket held by a single site never blocks the add, whatever its
// configs: the backfill creates it on the peers exactly like the
// historical one-non-empty-site path.
#[test]
fn test_validate_add_preflight_topology_ignores_unshared_bucket_configs() {
let local_peer = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let mut local = preflight_site("local", "https://local.example.com", "local-dep", 1);
let remote = preflight_site("remote", "https://remote.example.com", "remote-dep", 1);
local.buckets.insert(
"local-only".to_string(),
AddPreflightBucketCompat {
versioning_enabled: false,
object_lock_enabled: true,
},
);
let infos = vec![local, remote];
validate_add_preflight_topology(&infos, &local_peer).expect("unshared buckets should not block the add");
}
#[test]
@@ -9524,6 +9784,91 @@ mod tests {
assert!(err.to_string().contains("different site replication peer set"));
}
// add_preflight_bucket_compat reads the build_sr_info wire form:
// base64-encoded XML for both the versioning and the object-lock config.
#[test]
fn test_add_preflight_bucket_compat_parses_wire_configs() {
let info = SRBucketInfo {
versioning: Some(
BASE64_STANDARD.encode_to_string(b"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>"),
),
object_lock_config: Some(BASE64_STANDARD.encode_to_string(
b"<ObjectLockConfiguration><ObjectLockEnabled>Enabled</ObjectLockEnabled></ObjectLockConfiguration>",
)),
..Default::default()
};
let compat = add_preflight_bucket_compat("https://a.example.com", "b", &info).expect("wire configs should parse");
assert!(compat.versioning_enabled);
assert!(compat.object_lock_enabled);
}
#[test]
fn test_add_preflight_bucket_compat_absent_and_suspended_configs_are_disabled() {
let absent = add_preflight_bucket_compat("https://a.example.com", "b", &SRBucketInfo::default())
.expect("absent configs should parse");
assert!(!absent.versioning_enabled);
assert!(!absent.object_lock_enabled);
let suspended = SRBucketInfo {
versioning: Some(
BASE64_STANDARD
.encode_to_string(b"<VersioningConfiguration><Status>Suspended</Status></VersioningConfiguration>"),
),
..Default::default()
};
let compat =
add_preflight_bucket_compat("https://a.example.com", "b", &suspended).expect("suspended config should parse");
assert!(!compat.versioning_enabled, "suspended versioning is not merge-safe");
}
// rustfs/backlog#2070: a committed add must supersede this site's own
// half-finished removal (mirroring the join side, rustfs/rustfs#5963) —
// otherwise the reconcile tick replays the stale removal against the
// freshly re-paired peer and dismantles the new pairing.
#[test]
fn test_adopt_add_commit_state_clears_pending_remove() {
let mut state = SiteReplicationState {
pending_remove: Some(PendingRemove {
id: "remove-1".to_string(),
..Default::default()
}),
edit_generation: 7,
..Default::default()
};
let next_state = SiteReplicationState {
name: "local".to_string(),
peers: BTreeMap::from([
("local-dep".to_string(), peer("local", "https://local.example.com")),
("remote-dep".to_string(), peer("remote", "https://remote.example.com")),
]),
updated_at: Some(OffsetDateTime::now_utc()),
sync_state_initialized: true,
..Default::default()
};
adopt_add_commit_state(&mut state, next_state);
assert!(state.pending_remove.is_none(), "the committed add supersedes the removal");
assert_eq!(state.peers.len(), 2, "the add's topology is adopted");
assert_eq!(state.edit_generation, 7, "commit-owned fields keep the loaded value");
}
// Fail closed: a config this site cannot read must fail the preflight
// instead of defaulting into an unsafe admission.
#[test]
fn test_add_preflight_bucket_compat_rejects_undecodable_config() {
let info = SRBucketInfo {
versioning: Some(BASE64_STANDARD.encode_to_string(b"<VersioningConfiguration")),
..Default::default()
};
let err = add_preflight_bucket_compat("https://a.example.com", "b", &info).expect_err("broken XML should fail");
assert!(err.to_string().contains("unreadable versioning config"));
}
/// P1-15 review follow-up: the receiving side of the ordering fence. Two
/// nodes of the sending site can fan out in the opposite order to their
/// commits; the receiver decides ordering from the generation the sender
+15 -3
View File
@@ -20,8 +20,8 @@ use time::OffsetDateTime;
mod ecstore_bucket {
pub(crate) use crate::storage::storage_api::ecstore_bucket::{
bandwidth, bucket_target_sys, durability, lifecycle, metadata, metadata_sys, quota, replication, target, utils,
versioning, versioning_sys,
bandwidth, bucket_target_sys, durability, lifecycle, metadata, metadata_sys, object_lock, quota, replication, target,
utils, versioning, versioning_sys,
};
}
@@ -185,6 +185,16 @@ impl AdminVersioningConfigExt for s3s::dto::VersioningConfiguration {
}
}
pub(crate) trait AdminObjectLockConfigExt {
fn enabled(&self) -> bool;
}
impl AdminObjectLockConfigExt for s3s::dto::ObjectLockConfiguration {
fn enabled(&self) -> bool {
<s3s::dto::ObjectLockConfiguration as ecstore_bucket::object_lock::ObjectLockApi>::enabled(self)
}
}
pub(crate) mod bandwidth {
pub(crate) mod monitor {
pub(crate) type BandwidthDetails = super::super::ecstore_bucket::bandwidth::monitor::BandwidthDetails;
@@ -863,7 +873,9 @@ pub(crate) mod bucket {
pub(crate) use super::replication;
pub(crate) use super::target;
pub(crate) use super::versioning_sys;
pub(crate) use super::{AdminReplicationConfigExt, AdminVersioningConfigExt, is_reserved_or_invalid_bucket};
pub(crate) use super::{
AdminObjectLockConfigExt, AdminReplicationConfigExt, AdminVersioningConfigExt, is_reserved_or_invalid_bucket,
};
pub(crate) mod utils {
pub(crate) use super::super::ecstore_utils::{deserialize, is_valid_object_prefix, serialize};
File diff suppressed because it is too large Load Diff
+4
View File
@@ -182,6 +182,10 @@ fn object_s3_error(code: S3ErrorCode, message: impl Into<std::borrow::Cow<'stati
S3Error::with_message(code, message)
}
fn object_s3_error_default(code: S3ErrorCode) -> S3Error {
S3Error::new(code)
}
mod copy;
mod delete;
mod extract;
+6 -8
View File
@@ -308,14 +308,14 @@ pub(crate) fn guard_put_object_body_read_timeout(
})
}
struct PooledBufferReader {
pub(super) struct PooledBufferReader {
buffer: PooledBuffer,
len: usize,
pos: usize,
}
impl PooledBufferReader {
fn new(buffer: PooledBuffer, len: usize) -> Self {
pub(super) fn new(buffer: PooledBuffer, len: usize) -> Self {
Self { buffer, len, pos: 0 }
}
}
@@ -631,7 +631,7 @@ fn select_put_path_with_concurrency(
/// where the allocation cost is negligible (≤4KiB memcpy).
const POOL_BYPASS_MAX_SIZE: usize = 4 * 1024;
async fn read_small_put_body_into<R, B>(body: &mut R, buf: &mut B, size: usize) -> S3Result<()>
pub(super) async fn read_small_put_body_into<R, B>(body: &mut R, buf: &mut B, size: usize) -> S3Result<()>
where
R: AsyncRead + Unpin,
B: bytes::BufMut,
@@ -958,11 +958,9 @@ impl DefaultObjectUsecase {
return Err(s3_error!(InvalidStorageClass));
}
// An authorized inbound replication PUT must store the replica verbatim.
// A snowball-extracted member object keeps `x-amz-meta-snowball-auto-extract`
// in its user metadata, and the replication client replays stored metadata
// as headers re-dispatching that PUT into the extract path would try to
// untar the member's own bytes (failing replication for any non-archive
// member) instead of writing the replica.
// Legacy snowball-extracted members may still carry the auto-extract
// metadata, which replication replays as a header. Do not interpret that
// historical user metadata as a request to untar the member again.
let inbound_replication_put = replication_request_authorized(&req)
&& get_header(&req.headers, SUFFIX_SOURCE_REPLICATION_REQUEST).as_deref() == Some("true");
if max_content_length.is_some() && is_put_object_extract_requested(&req.headers) {
+7 -2
View File
@@ -978,9 +978,12 @@ pub(crate) mod bucket {
}
pub(crate) mod concurrency {
#[cfg(test)]
pub(crate) use crate::storage::storage_api::concurrency_consumer::SNOWBALL_MEMBER_COMMIT_LIMIT;
pub(crate) use crate::storage::storage_api::concurrency_consumer::{
ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, IoQueueStatus, IoStrategy,
PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
PutObjectGuard, SNOWBALL_STAGING_BYTES_LIMIT, get_concurrency_aware_buffer_size, get_concurrency_manager,
get_put_concurrency_aware_buffer_size,
};
}
@@ -1097,7 +1100,9 @@ pub(crate) mod s3_api {
}
pub(crate) mod tagging {
pub(crate) use crate::storage::storage_api::s3_api_consumer::tagging::resolve_copy_object_tags;
pub(crate) use crate::storage::storage_api::s3_api_consumer::tagging::{
parse_copy_object_tags, resolve_copy_object_tags,
};
}
}
+29 -1
View File
@@ -815,7 +815,15 @@ pub fn get_condition_values_with_query_and_client_info(
/// `key`, either because the server already derived that key from verified state or
/// because it is a well-known identity/context key that only the server may populate.
fn is_reserved_condition_key(key: &str, server_derived: &HashMap<String, Vec<String>>) -> bool {
server_derived.contains_key(key) || is_server_derived_condition_key(key)
server_derived.contains_key(key)
|| [
AMZ_OBJECT_LOCK_MODE_LOWER,
AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER,
AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER,
]
.iter()
.any(|header| key.eq_ignore_ascii_case(header.trim_start_matches("x-amz-")))
|| is_server_derived_condition_key(key)
}
/// Get request authentication type
@@ -1670,17 +1678,37 @@ mod tests {
let cred = create_test_credentials();
let mut headers = HeaderMap::new();
headers.insert(AMZ_OBJECT_LOCK_MODE_LOWER, HeaderValue::from_static("GOVERNANCE"));
headers.insert(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER, HeaderValue::from_static("OFF"));
headers.insert(AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER, HeaderValue::from_static("2024-12-31T23:59:59Z"));
headers.insert("object-lock-mode", HeaderValue::from_static("COMPLIANCE"));
headers.insert("object-lock-legal-hold", HeaderValue::from_static("ON"));
headers.insert("object-lock-retain-until-date", HeaderValue::from_static("2099-12-31T23:59:59Z"));
let conditions = get_condition_values(&headers, &cred, None, None, None);
assert_eq!(conditions.get("object-lock-mode"), Some(&vec!["GOVERNANCE".to_string()]));
assert_eq!(conditions.get("object-lock-legal-hold"), Some(&vec!["OFF".to_string()]));
assert_eq!(
conditions.get("object-lock-retain-until-date"),
Some(&vec!["2024-12-31T23:59:59Z".to_string()])
);
}
#[test]
fn object_lock_condition_aliases_cannot_spoof_canonical_headers() {
let cred = create_test_credentials();
let mut headers = HeaderMap::new();
headers.insert("object-lock-mode", HeaderValue::from_static("COMPLIANCE"));
headers.insert("object-lock-legal-hold", HeaderValue::from_static("ON"));
headers.insert("object-lock-retain-until-date", HeaderValue::from_static("2099-12-31T23:59:59Z"));
let conditions = get_condition_values(&headers, &cred, None, None, None);
assert_eq!(conditions.get("object-lock-mode"), None);
assert_eq!(conditions.get("object-lock-legal-hold"), None);
assert_eq!(conditions.get("object-lock-retain-until-date"), None);
}
#[test]
fn test_get_condition_values_with_grant_headers() {
let cred = create_test_credentials();
+171 -2
View File
@@ -31,10 +31,12 @@ use rustfs_io_metrics::bandwidth::{BandwidthMonitor, BandwidthSnapshot};
use rustfs_io_metrics::{MetricsCollector, PerformanceMetrics};
use std::sync::{Arc, LazyLock, Mutex};
use std::time::Duration;
use tokio::sync::Semaphore;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
use tracing::debug;
const DERIVED_LARGE_PUT_ADMISSION_LIMIT_MAX: usize = 32;
pub(crate) const SNOWBALL_MEMBER_COMMIT_LIMIT: usize = 32;
pub(crate) const SNOWBALL_STAGING_BYTES_LIMIT: usize = 4 * MI_B;
/// Global concurrency manager instance
pub(crate) static CONCURRENCY_MANAGER: LazyLock<ConcurrencyManager> = LazyLock::new(ConcurrencyManager::new);
@@ -69,6 +71,12 @@ pub struct ConcurrencyManager {
metrics_collector: Arc<MetricsCollector>,
/// Foreground write admission policy, resolved once at startup.
foreground_write_admission_policy: ForegroundWriteAdmissionPolicy,
/// Snowball members are internal PUTs, so they use a separate global gate
/// from preparation through the independently owned post-commit tail.
snowball_member_commit_semaphore: Arc<Semaphore>,
/// Bounds the owned member bodies and metadata retained between TAR parsing
/// and storage commit across all extract requests.
snowball_staging_bytes_semaphore: Arc<Semaphore>,
}
impl std::fmt::Debug for ConcurrencyManager {
@@ -417,6 +425,8 @@ impl ConcurrencyManager {
bandwidth_monitor,
metrics_collector,
foreground_write_admission_policy,
snowball_member_commit_semaphore: Arc::new(Semaphore::new(SNOWBALL_MEMBER_COMMIT_LIMIT)),
snowball_staging_bytes_semaphore: Arc::new(Semaphore::new(SNOWBALL_STAGING_BYTES_LIMIT)),
}
}
@@ -556,6 +566,40 @@ impl ConcurrencyManager {
.await
}
/// Admit a Snowball member through the foreground PUT policy using the
/// member's logical size. The outer archive has a separate preflight, while
/// every member shares the ordinary PUT gate and its wait/rejection policy.
pub(crate) async fn admit_snowball_foreground_write(
&self,
member_size: i64,
) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
self.foreground_write_admission_policy
.admit(ForegroundWriteAdmissionKind::PutObject, member_size)
.await
}
/// Acquire one global Snowball member lifecycle slot.
pub(crate) async fn acquire_snowball_member_commit(&self) -> Result<OwnedSemaphorePermit, tokio::sync::AcquireError> {
self.snowball_member_commit_semaphore.clone().acquire_owned().await
}
/// Try to acquire one global Snowball member lifecycle slot.
pub(crate) fn try_acquire_snowball_member_commit(&self) -> Option<OwnedSemaphorePermit> {
self.snowball_member_commit_semaphore.clone().try_acquire_owned().ok()
}
/// Try to reserve prepared-member bytes without waiting.
///
/// A producer holding a non-empty micro-batch must use this method and
/// flush before waiting, otherwise several archives can each retain part of
/// the global budget while waiting forever for the remainder.
pub(crate) fn try_acquire_snowball_staging_bytes(&self, bytes: u32) -> Option<OwnedSemaphorePermit> {
self.snowball_staging_bytes_semaphore
.clone()
.try_acquire_many_owned(bytes)
.ok()
}
/// Admit a multipart UploadPart request under the configured write gate.
///
/// Multipart workloads can saturate memory and internode write streams with
@@ -1050,13 +1094,138 @@ impl Default for ConcurrencyManager {
mod integration_tests {
use super::super::io_schedule::{IoLoadLevel, IoPriority};
use super::super::request_guard::GetObjectGuard;
use super::{ConcurrencyManager, ForegroundWriteAdmission, derive_large_put_admission_limit};
use super::{
ConcurrencyManager, ForegroundWriteAdmission, SNOWBALL_MEMBER_COMMIT_LIMIT, SNOWBALL_STAGING_BYTES_LIMIT,
derive_large_put_admission_limit,
};
use crate::storage::storage_api::concurrency_consumer::PutObjectGuard;
use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass};
use rustfs_io_core::io_profile::{AccessPattern, StorageMedia};
use serial_test::serial;
use std::time::Duration;
#[test]
fn test_snowball_gates_are_global_bounded_and_reusable() {
let manager = ConcurrencyManager::new();
let clone = manager.clone();
let commit_permits = manager
.snowball_member_commit_semaphore
.clone()
.try_acquire_many_owned(u32::try_from(SNOWBALL_MEMBER_COMMIT_LIMIT).expect("Snowball commit limit must fit into u32"))
.expect("the exact Snowball commit limit must be available");
assert!(
clone.snowball_member_commit_semaphore.clone().try_acquire_owned().is_err(),
"a cloned manager must share the global commit gate"
);
drop(commit_permits);
assert!(clone.snowball_member_commit_semaphore.clone().try_acquire_owned().is_ok());
let staging_bytes = u32::try_from(SNOWBALL_STAGING_BYTES_LIMIT).expect("Snowball staging limit must fit into u32");
let staging_permit = manager
.try_acquire_snowball_staging_bytes(staging_bytes)
.expect("the exact Snowball staging budget must be available");
assert!(
clone.try_acquire_snowball_staging_bytes(1).is_none(),
"a cloned manager must share the global staging budget"
);
drop(staging_permit);
assert!(clone.try_acquire_snowball_staging_bytes(1).is_some());
}
#[tokio::test]
async fn test_snowball_members_share_the_strict_foreground_put_gate() {
let manager = ConcurrencyManager::with_put_admission_for_test(true, 2, Duration::ZERO);
let outer = match manager
.admit_put_object(1)
.await
.expect("strict outer admission must remain open")
{
ForegroundWriteAdmission::Admitted(permit) => permit,
outcome => panic!("strict outer admission must return a permit: {outcome:?}"),
};
let member = match manager
.admit_snowball_foreground_write(1)
.await
.expect("strict member admission must remain open")
{
ForegroundWriteAdmission::Admitted(permit) => permit,
outcome => panic!("strict member admission must return a permit: {outcome:?}"),
};
assert!(
matches!(
manager
.admit_snowball_foreground_write(1)
.await
.expect("strict member admission must remain open"),
ForegroundWriteAdmission::Rejected
),
"a saturated Snowball member admission must preserve the zero-wait rejection policy"
);
assert!(
matches!(
manager.admit_put_object(1).await.expect("strict gate must remain usable"),
ForegroundWriteAdmission::Rejected
),
"outer PUTs and Snowball members must exhaust the same strict gate"
);
drop(outer);
let replacement = match manager
.admit_snowball_foreground_write(1)
.await
.expect("released strict capacity must be reusable")
{
ForegroundWriteAdmission::Admitted(permit) => permit,
outcome => panic!("strict replacement admission must return a permit: {outcome:?}"),
};
drop((member, replacement));
}
#[tokio::test]
async fn test_snowball_members_use_their_size_for_the_large_foreground_put_gate() {
let min_size = 16 * 1024 * 1024;
let manager = ConcurrencyManager::with_large_put_admission_for_test(true, 1, min_size, Duration::ZERO);
assert!(matches!(
manager
.admit_put_object((min_size - 1) as i64)
.await
.expect("small outer archive admission must remain open"),
ForegroundWriteAdmission::Disabled
));
let large_member = match manager
.admit_snowball_foreground_write(min_size as i64)
.await
.expect("large Snowball member admission must remain open")
{
ForegroundWriteAdmission::Admitted(permit) => permit,
outcome => panic!("large Snowball member must consume the large PUT gate: {outcome:?}"),
};
assert!(matches!(
manager
.admit_snowball_foreground_write(min_size as i64)
.await
.expect("saturated Snowball member admission must remain open"),
ForegroundWriteAdmission::Rejected
));
assert!(matches!(
manager
.admit_put_object(min_size as i64)
.await
.expect("ordinary large PUT admission must remain open"),
ForegroundWriteAdmission::Rejected
));
assert!(matches!(
manager
.admit_snowball_foreground_write((min_size - 1) as i64)
.await
.expect("small Snowball member admission must remain open"),
ForegroundWriteAdmission::Disabled
));
drop(large_member);
}
#[tokio::test]
#[serial]
async fn test_concurrency_manager_priority_queue_integration() {
+3
View File
@@ -51,6 +51,9 @@ pub use io_schedule::{
pub use request_guard::{GetObjectGuard, PutObjectGuard};
// Concurrency manager
#[cfg(test)]
pub(crate) use manager::SNOWBALL_MEMBER_COMMIT_LIMIT;
pub(crate) use manager::SNOWBALL_STAGING_BYTES_LIMIT;
pub use manager::{ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission};
// ============================================
+5 -2
View File
@@ -125,9 +125,12 @@ pub(crate) mod access_consumer {
}
pub(crate) mod concurrency_consumer {
#[cfg(test)]
pub(crate) use super::super::concurrency::SNOWBALL_MEMBER_COMMIT_LIMIT;
pub(crate) use super::super::concurrency::{
ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, IoQueueStatus, IoStrategy,
PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
PutObjectGuard, SNOWBALL_STAGING_BYTES_LIMIT, get_concurrency_aware_buffer_size, get_concurrency_manager,
get_put_concurrency_aware_buffer_size,
};
}
@@ -350,7 +353,7 @@ pub(crate) mod s3_api_consumer {
}
pub(crate) mod tagging {
pub(crate) use super::super::super::s3_api::tagging::resolve_copy_object_tags;
pub(crate) use super::super::super::s3_api::tagging::{parse_copy_object_tags, resolve_copy_object_tags};
}
}
+1 -1
View File
@@ -16445,7 +16445,7 @@ fn object_mutation_entrypoints_call_reserved_prefix_guard() {
"if let Err(err) = validate_table_catalog_object_mutation(&bucket, &obj_id.key).await",
"validate_table_catalog_object_mutation(&bucket, &object).await?;",
"validate_object_key(&key, \"PUT\")?;\n validate_table_catalog_object_mutation(&bucket, &key).await?;",
"validate_table_catalog_object_mutation(&bucket, &fpath).await?;",
"extract_try!(validate_table_catalog_object_mutation(&bucket, &fpath).await);",
] {
assert!(source.contains(expected), "missing object mutation guard: {expected}");
}