Compare commits

..

6 Commits

Author SHA1 Message Date
houseme 2600a50c5d Merge branch 'main' into overtrue/fix-read-version-quorum-fixture 2026-07-29 09:57:33 +08:00
cxymds c1538cf1c3 fix(multipart): serialize complete and abort (#5356)
* fix(multipart): serialize complete and abort

* test(multipart): order abort-first finalization

* fix(multipart): enforce quorum staging cleanup

* fix(multipart): remove stale mutable binding
2026-07-29 01:49:19 +00:00
houseme 5af56cbb02 test(ci): stabilize lifecycle timeout coverage (#5404)
Co-authored-by: heihutu <heihutu@gmail.com>
2026-07-29 01:37:29 +00:00
houseme f99956eade test(lifecycle): classify mixed rollout harness results (#5384)
* test(lifecycle): classify mixed rollout harness results

Classify Docker #1508 evidence as strict, baseline, blocked, or failed so tiered-storage baseline runs cannot be mistaken for strict mixed-version rollout closure evidence.

Co-Authored-By: heihutu <heihutu@gmail.com>

* test: classify Docker manual transition preemption

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-07-29 01:21:23 +00:00
houseme 775279b6fd fix(notify): defer disabled bootstrap storage refresh (#5373)
Co-authored-by: heihutu <heihutu@gmail.com>
2026-07-28 17:39:22 +00:00
overtrue 10db415099 test(ecstore): give the read-version quorum fixture a valid part
`read_version_quorum_test_set` (#5365) writes `valid_test_fileinfo` to disk
before exercising the optimized version read. That fixture deliberately
carries a positive size with no parts so it can drive
`get_object_with_fileinfo_rejects_positive_size_without_parts`, and #5354
subsequently taught `validate_collection_contents` to reject exactly that
shape as `FileCorrupt`. #5354 fixed up the fixtures that existed when it
landed, but #5365 was developed in parallel and its new fixture was not,
so `write_metadata` now fails in the helper and both
`read_version_optimized_counts_only_valid_metadata_toward_quorum` and
`read_version_optimized_uses_current_set_drive_quorum_not_pool_set_count`
panic on `metadata should be written before quorum read: FileCorrupt`.

Neither PR could see the break on its own; it only appears once both are on
main, where it fails the Test and Lint lane deterministically for every PR.

Push a part matching the fixture's size in the helper — the same fix #5354
applied to its own on-disk fixture in `disk/local.rs` — instead of changing
`valid_test_fileinfo`, which must stay part-less for the rejection test.

Verification:
- cargo test -p rustfs-ecstore --lib set_disk::read::metadata_cache_tests
  (previously 2 failed, now 33 passed / 0 failed)
- cargo fmt --all --check and cargo clippy -p rustfs-ecstore --lib --tests clean
2026-07-29 00:21:27 +08:00
43 changed files with 2824 additions and 1471 deletions
+9 -14
View File
@@ -1,17 +1,14 @@
# nextest configuration for RustFS.
#
# Serialize two known load-sensitive / global-state-sharing ecstore test groups
# so the full parallel nextest suite stops producing spurious failures
# (backlog #937). These tests pass in isolation but flake under the loaded
# parallel run for two distinct reasons:
# Serialize the ecstore tests that share the process-wide disk registry or
# exercise a multi-disk commit handoff across nextest process boundaries.
#
# * store::bucket::tests::bucket_delete_* share process/global state (disk
# registry, lock client) and race make_bucket into InsufficientWriteQuorum
# when run concurrently with other ecstore tests.
# * bucket_lifecycle_ops::tests::concurrent_resend_same_part_commits_one_generation
# asserts a lock-acquire correctness property whose serialized cross-disk
# commits exceed the (already max'd, 60s) acquire deadline only when the
# suite saturates disk I/O.
# uses the shared multipart fixture and a deterministic uploadId-lock
# handoff, so it must not overlap another process mutating that fixture.
#
# serial_test's #[serial] attribute does NOT serialize these across runs:
# nextest executes each test in its own process, where the in-process
@@ -100,13 +97,6 @@ path = "junit.xml"
# profile's own overrides list, not the default profile's).
# ===========================================================================
# QUARANTINE: OPEN backlog#937 — concurrent_resend lock-acquire deadline flakes
# under saturated disk I/O in the full parallel suite.
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & test(concurrent_resend_same_part_commits_one_generation)'
test-group = 'ecstore-serial-flaky'
retries = 2
# QUARANTINE: OPEN backlog#937 — store::bucket::tests::bucket_delete_* race
# make_bucket into InsufficientWriteQuorum via shared global state under load.
[[profile.ci.overrides]]
@@ -114,6 +104,11 @@ filter = 'package(rustfs-ecstore) & test(/^store::bucket::tests::bucket_delete_(
test-group = 'ecstore-serial-flaky'
retries = 2
# Keep the deterministic multipart handoff isolated across nextest processes.
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & test(concurrent_resend_same_part_commits_one_generation)'
test-group = 'ecstore-serial-flaky'
# QUARANTINE: OPEN rustfs#4690 — walk_dir stall-budget accounting test depends
# on producer/consumer timing windows that stretch past the budget on loaded
# CI runners (regression test for rustfs#4644; failed on a zero-Rust-diff PR).
+1 -1
View File
@@ -141,7 +141,7 @@ jobs:
name: Test and Lint
if: github.event_name != 'pull_request' || github.event.action != 'closed'
runs-on: sm-standard-4
timeout-minutes: 60
timeout-minutes: 90
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
steps:
Generated
+31 -29
View File
@@ -1703,9 +1703,9 @@ dependencies = [
[[package]]
name = "camino"
version = "1.2.4"
version = "1.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5f2d30e4173c4026932d51d31d6b0613b1fd3014bf3f9f8943d4ba139c437ba0"
checksum = "bb1307f12aa967b5a58416e87b3653360e0fd614a016b6e970db08fecbb1b80d"
dependencies = [
"serde_core",
]
@@ -3625,13 +3625,13 @@ dependencies = [
[[package]]
name = "displaydoc"
version = "0.2.6"
version = "0.2.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1ac70aa55017e108007fbaf5aa0f54b021c98f92ff8af59d42eda9da96e3dd4f"
checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.119",
"syn 3.0.3",
]
[[package]]
@@ -3941,16 +3941,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
name = "event-listener"
version = "5.4.1"
version = "5.4.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab"
checksum = "5a23add41df1562121a9393cb065eab5146a1242410f23a644851e90cfd669d2"
dependencies = [
"concurrent-queue",
"parking",
"pin-project-lite",
]
@@ -5361,7 +5360,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46"
dependencies = [
"hermit-abi",
"libc",
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -6110,9 +6109,9 @@ dependencies = [
[[package]]
name = "metrique"
version = "0.1.28"
version = "0.1.29"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "aa466af30a9fe0b1db1dae097e0ce7ddac4b666b1d3a9f5c682e889272511dd8"
checksum = "d2e394c63e2d1a30aeb3b9392ecf3439d8475d2df810a8f4f6e66d6866754017"
dependencies = [
"itoa",
"jiff",
@@ -6140,9 +6139,9 @@ dependencies = [
[[package]]
name = "metrique-macro"
version = "0.1.19"
version = "0.1.20"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1f5febfaf14fea234b60e0ad43c728db274033893a38d0fa1f87d9b7056d3d61"
checksum = "786df1fd0abebd0db685f7e9a353c78756d4b370fb98a52376c2015fa55f141f"
dependencies = [
"Inflector",
"darling 0.23.0",
@@ -6169,9 +6168,9 @@ checksum = "2faca4e4480069ff02b1763b3b79f5cec7e8628e24d9dc5b6073f53d2577a4d9"
[[package]]
name = "metrique-writer"
version = "0.1.24"
version = "0.1.25"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "124326a2ac4c4f61562fa4d071735a1c463f9ba0317d1564f75dc01313dc12d7"
checksum = "82cdde44d241dab7fc8b7a32e0eb5dae6cd28f8de80b59f9a1e9f2f0b05e485e"
dependencies = [
"ahash",
"crossbeam-queue",
@@ -6190,9 +6189,9 @@ dependencies = [
[[package]]
name = "metrique-writer-core"
version = "0.1.18"
version = "0.1.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "55b5bbb6d88bde29f6ed74a574cbe49e51ea0c36cccc7e63eee2040db5675b15"
checksum = "e57379b7ee2272efaeaaa6de062503563e57333b24aadc7f2255b3d602899e8b"
dependencies = [
"derive-where",
"itertools 0.14.0",
@@ -8097,7 +8096,7 @@ dependencies = [
"once_cell",
"socket2",
"tracing",
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -9106,6 +9105,7 @@ dependencies = [
name = "rustfs-ecstore"
version = "1.0.0-beta.11"
dependencies = [
"aes-gcm",
"arc-swap",
"async-channel",
"async-recursion",
@@ -9121,6 +9121,7 @@ dependencies = [
"byteorder",
"bytes",
"bytesize",
"chacha20poly1305",
"chrono",
"criterion",
"enumset",
@@ -9173,6 +9174,7 @@ dependencies = [
"rustfs-erasure-codec",
"rustfs-filemeta",
"rustfs-io-metrics",
"rustfs-kms",
"rustfs-lifecycle",
"rustfs-lock",
"rustfs-madmin",
@@ -10222,7 +10224,7 @@ dependencies = [
"errno",
"libc",
"linux-raw-sys",
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -10295,7 +10297,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -10458,9 +10460,9 @@ dependencies = [
[[package]]
name = "schemars"
version = "1.2.1"
version = "1.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a2b42f36aa1cd011945615b92222f6bf73c599a102a300334cd7f8dbeec726cc"
checksum = "687274d293b6cdc6e73e0fee520bf2049650090d7164f87672d212a3c530cf4a"
dependencies = [
"dyn-clone",
"ref-cast",
@@ -10693,7 +10695,7 @@ dependencies = [
"indexmap 1.9.3",
"indexmap 2.14.0",
"schemars 0.9.0",
"schemars 1.2.1",
"schemars 1.2.2",
"serde_core",
"serde_json",
"serde_with_macros",
@@ -11495,10 +11497,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd"
dependencies = [
"fastrand",
"getrandom 0.3.4",
"getrandom 0.4.3",
"once_cell",
"rustix",
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -11873,9 +11875,9 @@ dependencies = [
[[package]]
name = "toml_parser"
version = "1.1.2+spec-1.1.0"
version = "1.1.3+spec-1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a2abe9b86193656635d2411dc43050282ca48aa31c2451210f4202550afb7526"
checksum = "1d38ac1cf9b95face32296c0a3ede1fdc270627c9d9c02a7274dd6d960dc4d56"
dependencies = [
"winnow",
]
@@ -12599,7 +12601,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.52.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -45,10 +45,10 @@
//! * Parity reconstruction: one data disk is taken offline
//! (`take_disk_offline`) and the SAME object matrix is GET both ways while
//! the EC 2+2 set rebuilds each large object from the surviving shards. The
//! codec-streaming reader gate never inspects drive health, so the codec
//! fast path is exercised end-to-end through reconstruction; the test
//! asserts byte- and header-equality vs the legacy path AND that the codec
//! phase never fell back to a duplex pipe while reconstructing.
//! eager first/single-part setup may keep its conservative whole-request
//! fallback when shard placement makes codec streaming unsafe, so this phase
//! asserts byte- and header-equality vs the legacy path rather than requiring
//! zero duplex fallbacks under degraded drive health.
//! * Missing object: a GET for an absent key is compared across both phases
//! to prove the error semantics (HTTP status + S3 error code) are identical
//! — the codec env must not perturb the NoSuchKey negative path.
@@ -475,15 +475,14 @@ mod tests {
"ranged GET length diverged with codec streaming enabled"
);
// ---- Phase B degraded: the same reconstruction, now on the codec path ----
// Re-run the reconstruction A/B with the codec-streaming gates still
// open. The reader gate decision is independent of drive health (it
// never inspects disk state), so the codec fast path is exercised
// end-to-end while the EC set rebuilds each large object from the
// surviving shards — this is a real codec-vs-legacy reconstruction test,
// not legacy-vs-legacy. Snapshot the duplex count first (the range GET
// above already used the duplex path) so we can measure only the markers
// these degraded codec GETs add.
// ---- Phase B degraded: the same reconstruction, with codec gates open ----
// Re-run the reconstruction A/B with codec-streaming enabled. If eager
// first/single-part setup cannot prove the codec path is safe for the
// surviving shards, the implementation intentionally preserves the
// whole-request legacy fallback; later multipart parts can degrade in
// place. This phase verifies parity-reconstructed bytes and headers,
// while the healthy phase above remains the strict zero-duplex path
// confirmation.
let dup_codec_before_degraded = count_marker(&codec_log, DUPLEX_MARKER);
harness.take_disk_offline(0)?;
let mut codec_degraded: BTreeMap<String, GetView> = BTreeMap::new();
@@ -511,16 +510,11 @@ mod tests {
);
}
// Path confirmation under reconstruction: the codec fast path must have
// served the reconstructed large objects without ever falling back to
// the legacy duplex pipe. Without this, the equivalence above could be
// legacy-vs-legacy and prove nothing about codec reconstruction.
// Keep degraded duplex markers as diagnostic evidence only: eager setup
// may fall back before streaming when shard safety cannot be proven.
sleep(Duration::from_millis(300)).await;
let dup_codec_degraded = count_marker(&codec_log, DUPLEX_MARKER).saturating_sub(dup_codec_before_degraded);
assert_eq!(
dup_codec_degraded, 0,
"codec phase created {dup_codec_degraded} duplex pipe(s) while reconstructing large objects with disk0 offline; the codec fast path was not exercised under degraded reads (see {codec_log})"
);
info!(dup_codec_degraded, "codec phase degraded-read legacy duplex marker count");
info!(
objects = baseline.len(),
+3
View File
@@ -57,6 +57,7 @@ rustfs-policy.workspace = true
rustfs-protos.workspace = true
rustfs-replication.workspace = true
rustfs-lifecycle.workspace = true
rustfs-kms.workspace = true
rustfs-s3-types = { workspace = true }
rustfs-data-usage.workspace = true
rustfs-object-capacity.workspace = true
@@ -123,6 +124,8 @@ libc.workspace = true
rustix = { workspace = true, features = ["process", "fs"] }
rustfs-madmin.workspace = true
reqwest = { workspace = true }
aes-gcm = { workspace = true, features = ["rand_core"] }
chacha20poly1305.workspace = true
aws-sdk-s3 = { workspace = true, default-features = false, features = ["sigv4a", "default-https-client", "rt-tokio"] }
urlencoding = { workspace = true }
smallvec = { workspace = true, features = ["serde"] }
+4 -6
View File
@@ -381,12 +381,10 @@ pub mod notification {
pub mod object {
pub use crate::object_api::{
BLOCK_SIZE_V2, ERASURE_ALGORITHM, EncryptionResolutionError, EncryptionResolutionErrorKind, GetObjectBodyCacheHook,
GetObjectBodyCacheHookLookup, GetObjectBodySource, GetObjectReader, ObjectEncryptionResolver, ObjectInfo,
ObjectMutationHook, ObjectOptions, PutObjReader, RangedDecompressReader, ReadEncryptionMaterial, ReadEncryptionMode,
ReadEncryptionRequest, StreamConsumer, get_object_body_cache_plaintext_len, lookup_get_object_body_cache_hook,
register_get_object_body_cache_hook, register_object_mutation_hook, unregister_get_object_body_cache_hook,
unregister_object_mutation_hook,
BLOCK_SIZE_V2, ERASURE_ALGORITHM, GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, GetObjectBodySource,
GetObjectReader, ObjectInfo, ObjectMutationHook, ObjectOptions, PutObjReader, RangedDecompressReader, StreamConsumer,
get_object_body_cache_plaintext_len, lookup_get_object_body_cache_hook, register_get_object_body_cache_hook,
register_object_mutation_hook, unregister_get_object_body_cache_hook, unregister_object_mutation_hook,
};
pub use crate::store::PreparedGetObjectReader;
}
@@ -11112,6 +11112,7 @@ mod tests {
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn concurrent_resend_same_part_commits_one_generation() {
use crate::set_disk::{MultipartCommitBarrier, MultipartCommitPause};
use crate::storage_api_contracts::object::ObjectIO as _;
let (_paths, ecstore) = setup_test_env().await;
@@ -11133,51 +11134,36 @@ mod tests {
})
.collect();
// Two independent causes can produce a spurious lock-acquire timeout
// here, and both must stay covered:
// 1. A lost/stolen fast-lock wakeup could strand a waiter until the
// deadline — fixed for real in fast_lock::shard by bounding each
// notification wait (NOTIFY_WAIT_CAP re-polling).
// 2. Under the full nextest suite on loaded CI disks, the
// *legitimately serialized* cross-disk commits can exceed the
// acquire deadline all by themselves — observed on CI at the 5s
// default and the 30s production default with six resends, and
// again at 60s, which is a hard ceiling: fast_lock clamps every
// requested timeout to MAX_ACQUIRE_TIMEOUT (60s), so raising the
// env override higher is a no-op (the Timeout error still reports
// the requested value). Keep the guard about the correctness
// property, not disk latency: request the full 60s ceiling and cap
// the queue depth at three resends, so the last waiter sits behind
// at most two serialized commits (~12s each on the slowest observed
// CI runner, comfortably inside the deadline). Three concurrent
// resends still race the streaming phase and contend on the commit
// lock, which is all the generation-mixing regression needs.
// `#[serial]` keeps the process-wide env override isolated.
let results = temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("60"))], async {
let mut tasks = tokio::task::JoinSet::new();
for payload in candidates.iter().cloned() {
let store = ecstore.clone();
let bucket = bucket.clone();
let upload_id = upload.upload_id.clone();
tasks.spawn(async move {
let mut data = PutObjReader::from_vec(payload.clone());
store
.put_object_part(&bucket, object, &upload_id, 1, &mut data, &ObjectOptions::default())
.await
.map(|info| (info, payload))
});
}
let commit_barrier = MultipartCommitBarrier::install(&bucket, object, MultipartCommitPause::PutPartBeforeLockLost);
let start = Arc::new(tokio::sync::Barrier::new(candidates.len() + 1));
let mut tasks = tokio::task::JoinSet::new();
for payload in candidates.iter().cloned() {
let store = ecstore.clone();
let bucket = bucket.clone();
let upload_id = upload.upload_id.clone();
let start = Arc::clone(&start);
tasks.spawn(async move {
start.wait().await;
let mut data = PutObjReader::from_vec(payload.clone());
store
.put_object_part(&bucket, object, &upload_id, 1, &mut data, &ObjectOptions::default())
.await
.map(|info| (info, payload))
});
}
start.wait().await;
// Every concurrent resend must succeed; the commit lock must never
// starve a waiter into a timeout.
let mut results = Vec::new();
while let Some(joined) = tasks.join_next().await {
let outcome = joined.expect("put_object_part task should not panic");
results.push(outcome.expect("every concurrent same-part resend must succeed without lock timeout"));
}
results
})
.await;
// The first writer holds the uploadId commit lock while the other
// resends reach the same critical section. Releasing it proves the
// handoff without depending on saturated CI disk latency.
commit_barrier.wait_until_paused().await;
commit_barrier.release();
let mut results = Vec::new();
while let Some(joined) = tasks.join_next().await {
let outcome = joined.expect("put_object_part task should not panic");
results.push(outcome.expect("every concurrent same-part resend must succeed without lock timeout"));
}
assert_eq!(results.len(), candidates.len());
// Exactly one generation is visible after the serialized commits, and its
+15 -1
View File
@@ -46,6 +46,16 @@ lazy_static! {
m.insert("x-amz-replication-status".to_string(), true);
m
};
static ref SSE_HEADERS: HashMap<String, bool> = {
let mut m = HashMap::new();
m.insert("x-amz-server-side-encryption".to_string(), true);
m.insert("x-amz-server-side-encryption-aws-kms-key-id".to_string(), true);
m.insert("x-amz-server-side-encryption-context".to_string(), true);
m.insert("x-amz-server-side-encryption-customer-algorithm".to_string(), true);
m.insert("x-amz-server-side-encryption-customer-key".to_string(), true);
m.insert("x-amz-server-side-encryption-customer-key-md5".to_string(), true);
m
};
}
pub fn is_standard_query_value(qs_key: &str) -> bool {
@@ -60,12 +70,16 @@ pub fn is_standard_header(header_key: &str) -> bool {
*SUPPORTED_HEADERS.get(&header_key.to_lowercase()).unwrap_or(&false)
}
pub fn is_sse_header(header_key: &str) -> bool {
*SSE_HEADERS.get(&header_key.to_lowercase()).unwrap_or(&false)
}
pub fn is_amz_header(header_key: &str) -> bool {
let key = header_key.to_lowercase();
key.starts_with("x-amz-meta-")
|| key.starts_with("x-amz-grant-")
|| key == "x-amz-acl"
|| rustfs_utils::http::is_sse_header(header_key)
|| is_sse_header(header_key)
|| key.starts_with("x-amz-checksum-")
}
@@ -1,80 +0,0 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use async_trait::async_trait;
use http::{HeaderMap, HeaderValue};
use std::collections::HashMap;
use std::error::Error;
use std::fmt::{Display, Formatter};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReadEncryptionMode {
Direct { base_nonce: [u8; 12] },
Object,
}
pub struct ReadEncryptionMaterial {
pub key_bytes: [u8; 32],
pub mode: ReadEncryptionMode,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EncryptionResolutionErrorKind {
InvalidRequest,
InvalidMetadata,
ServiceUnavailable,
DecryptionFailed,
}
#[derive(Debug)]
pub struct EncryptionResolutionError {
kind: EncryptionResolutionErrorKind,
message: String,
}
impl EncryptionResolutionError {
pub fn new(kind: EncryptionResolutionErrorKind, message: impl Into<String>) -> Self {
Self {
kind,
message: message.into(),
}
}
pub fn kind(&self) -> EncryptionResolutionErrorKind {
self.kind
}
}
impl Display for EncryptionResolutionError {
fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
formatter.write_str(&self.message)
}
}
impl Error for EncryptionResolutionError {}
pub struct ReadEncryptionRequest<'a> {
pub bucket: &'a str,
pub object: &'a str,
pub metadata: &'a HashMap<String, String>,
pub headers: &'a HeaderMap<HeaderValue>,
}
#[async_trait]
pub trait ObjectEncryptionResolver: Send + Sync {
async fn resolve_read_material(
&self,
request: ReadEncryptionRequest<'_>,
) -> Result<Option<ReadEncryptionMaterial>, EncryptionResolutionError>;
}
-5
View File
@@ -84,7 +84,6 @@ pub(crate) fn legacy_encrypted_range_seek_enabled() -> bool {
}
mod body_cache_hook;
mod encryption;
mod hook_slot;
mod object_mutation_hook;
mod readers;
@@ -99,10 +98,6 @@ pub use body_cache_hook::{
pub(crate) use body_cache_hook::{
get_object_body_cache_hook, get_object_body_cache_hook_suppressed, without_get_object_body_cache_hook,
};
pub use encryption::{
EncryptionResolutionError, EncryptionResolutionErrorKind, ObjectEncryptionResolver, ReadEncryptionMaterial,
ReadEncryptionMode, ReadEncryptionRequest,
};
pub(crate) use object_mutation_hook::notify_object_mutation;
pub use object_mutation_hook::{ObjectMutationHook, register_object_mutation_hook, unregister_object_mutation_hook};
pub use readers::*;
File diff suppressed because it is too large Load Diff
+46 -17
View File
@@ -271,9 +271,29 @@ impl ObjectInfo {
}
pub fn is_encrypted(&self) -> bool {
self.user_defined
.keys()
.any(|key| rustfs_utils::http::is_object_encryption_marker(key))
// Corresponding to the logic in rustfs/src/sse.rs/encryption_material_to_metadata function
use rustfs_utils::http::{SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER};
self.user_defined.keys().any(|key| {
let lower = key.to_ascii_lowercase();
lower.starts_with("x-minio-encryption-")
|| lower.starts_with("x-minio-internal-server-side-encryption-")
|| matches!(
lower.as_str(),
"x-minio-internal-encrypted-multipart"
| "x-rustfs-encryption-key"
| "x-rustfs-encryption-algorithm"
| "x-rustfs-encryption-iv"
| "x-rustfs-encryption-key-id"
| "x-rustfs-encryption-context"
| "x-rustfs-encryption-tag"
| "x-amz-server-side-encryption-aws-kms-key-id"
| SSEC_ALGORITHM_HEADER
| SSEC_KEY_HEADER
| SSEC_KEY_MD5_HEADER
| "x-amz-server-side-encryption"
)
})
}
/// Maximum inline size for non-versioned objects (128 KiB).
@@ -317,7 +337,26 @@ impl ObjectInfo {
}
pub fn encryption_original_size(&self) -> std::io::Result<Option<i64>> {
rustfs_utils::http::get_object_encryption_original_size(&self.user_defined)
let actual_size = rustfs_utils::http::get_str(&self.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE);
if let Some(size_str) = self
.user_defined
.get("x-rustfs-encryption-original-size")
.map(String::as_str)
.or_else(|| {
self.user_defined
.get("x-amz-server-side-encryption-customer-original-size")
.map(String::as_str)
})
.or(actual_size.as_deref())
&& !size_str.is_empty()
{
let size = size_str
.parse::<i64>()
.map_err(|e| std::io::Error::other(format!("Failed to parse encryption original size: {e}")))?;
return Ok(Some(size));
}
Ok(None)
}
pub fn decrypted_size(&self) -> std::io::Result<i64> {
@@ -347,6 +386,9 @@ impl ObjectInfo {
return Ok(actual_size);
}
// Check if object is encrypted
// Managed SSE stores original size in x-rustfs-encryption-original-size metadata
// SSE-C stores original size in x-amz-server-side-encryption-customer-original-size
if let Some(size) = self.encryption_original_size()? {
return Ok(size);
}
@@ -836,19 +878,6 @@ mod tests {
assert!(!object.is_inline_fast_path_eligible(), "transitioned objects must fall back");
}
#[test]
fn minio_internal_encryption_metadata_is_not_treated_as_plaintext() {
let object = ObjectInfo {
user_defined: Arc::new(HashMap::from([(
"X-Minio-Internal-Server-Side-Encryption-Sealed-Key".to_string(),
"sealed".to_string(),
)])),
..Default::default()
};
assert!(object.is_encrypted());
}
#[test]
fn versions_after_marker_handles_null_version_marker() {
let first_version = Uuid::parse_str("11111111-2222-3333-4444-555555555555").unwrap();
-17
View File
@@ -46,7 +46,6 @@ use crate::bucket::metadata_sys::BucketMetadataSys;
use crate::bucket::replication::{DynReplicationPool, ReplicationStats};
use crate::disk::DiskStore;
use crate::layout::endpoints::{EndpointServerPools, SetupType};
use crate::object_api::ObjectEncryptionResolver;
use crate::services::event_notification::EventNotifier;
use crate::services::tier::tier::TierConfigMgr;
use rustfs_lock::{GlobalLockManager, get_global_lock_manager};
@@ -160,8 +159,6 @@ pub struct InstanceContext {
/// workers (scanner/heal/tier/lifecycle) without touching another instance.
/// Replaces the process-global cancel-token static.
background_cancel_token: OnceLock<CancellationToken>,
/// Resolves object-encryption material at the application boundary.
object_encryption_resolver: OnceLock<Arc<dyn ObjectEncryptionResolver>>,
tier_delete_journal_recovery_stores: std::sync::Mutex<HashSet<Uuid>>,
transition_transaction_recovery_stores: std::sync::Mutex<HashSet<Uuid>>,
#[cfg(test)]
@@ -200,7 +197,6 @@ impl InstanceContext {
local_disk_set_drives: Arc::new(RwLock::new(Vec::new())),
bucket_metadata_sys: std::sync::Mutex::new(None),
background_cancel_token: OnceLock::new(),
object_encryption_resolver: OnceLock::new(),
tier_delete_journal_recovery_stores: std::sync::Mutex::new(HashSet::new()),
transition_transaction_recovery_stores: std::sync::Mutex::new(HashSet::new()),
#[cfg(test)]
@@ -213,19 +209,6 @@ impl InstanceContext {
self.lock_manager.clone()
}
/// Install the application-owned object-encryption resolver once.
pub fn set_object_encryption_resolver(
&self,
resolver: Arc<dyn ObjectEncryptionResolver>,
) -> Result<(), Arc<dyn ObjectEncryptionResolver>> {
self.object_encryption_resolver.set(resolver)
}
/// Return the configured object-encryption resolver, if startup installed one.
pub fn object_encryption_resolver(&self) -> Option<&dyn ObjectEncryptionResolver> {
self.object_encryption_resolver.get().map(Arc::as_ref)
}
/// Set this instance's S3 region.
///
/// Write-once: panics on a second write, preserving the startup fail-fast
+5
View File
@@ -46,6 +46,7 @@ use crate::{
use rustfs_concurrency::WorkloadAdmissionSnapshotProvider;
use rustfs_config::server_config::{Config, get_global_server_config, set_global_server_config};
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_kms::{ObjectEncryptionService, get_global_encryption_service};
use rustfs_lock::client::LockClient;
use s3s::dto::BucketLifecycleConfiguration;
use s3s::region::Region;
@@ -104,6 +105,10 @@ pub(crate) fn record_erasure_write_quorum_failure(stage: &'static str, dominant_
global_internode_metrics().record_erasure_write_quorum_failure(stage, dominant_error);
}
pub(crate) async fn object_encryption_service() -> Option<Arc<ObjectEncryptionService>> {
get_global_encryption_service().await
}
pub fn object_store_handle() -> Option<Arc<ECStore>> {
resolve_object_store_handle()
}
+3 -1
View File
@@ -505,7 +505,9 @@ impl SetDisks {
}
fn file_info_has_encryption_metadata(meta: &FileInfo) -> bool {
meta.metadata.keys().any(|name| http::is_object_encryption_marker(name))
meta.metadata
.keys()
.any(|name| http::is_encryption_metadata_key(name) || http::is_sse_header(name))
}
fn starts_with_ignore_ascii_case(value: &str, prefix: &str) -> bool {
+9 -6
View File
@@ -143,17 +143,15 @@ use rustfs_object_capacity::capacity_scope::{
CapacityScope, CapacityScopeDisk, current_dirty_generation, record_capacity_scope, record_global_dirty_scope,
};
use rustfs_s3_types::EventName;
#[cfg(test)]
use rustfs_utils::http::SSEC_ALGORITHM_HEADER;
use rustfs_utils::http::headers::AMZ_OBJECT_TAGGING;
use rustfs_utils::http::headers::AMZ_STORAGE_CLASS;
use rustfs_utils::http::headers::{
CACHE_CONTROL, CONTENT_DISPOSITION, CONTENT_ENCODING, CONTENT_LANGUAGE, CONTENT_TYPE, EXPIRES, HeaderExt as _,
};
use rustfs_utils::http::{
SUFFIX_ACTUAL_OBJECT_SIZE_CAP, SUFFIX_ACTUAL_SIZE, SUFFIX_COMPRESSION, SUFFIX_COMPRESSION_SIZE, SUFFIX_REPLICATION_SSEC_CRC,
SUFFIX_RESTORE_OPERATION_ID, contains_key_str, get_header_map, get_str, insert_str, is_object_encryption_marker,
remove_header_map,
SSEC_ALGORITHM_HEADER, SSEC_KEY_HEADER, SSEC_KEY_MD5_HEADER, SUFFIX_ACTUAL_OBJECT_SIZE_CAP, SUFFIX_ACTUAL_SIZE,
SUFFIX_COMPRESSION, SUFFIX_COMPRESSION_SIZE, SUFFIX_REPLICATION_SSEC_CRC, SUFFIX_RESTORE_OPERATION_ID, contains_key_str,
get_header_map, get_str, insert_str, is_encryption_metadata_key, remove_header_map,
};
use rustfs_utils::{
HashAlgorithm,
@@ -409,7 +407,10 @@ pub(crate) fn strip_internal_multipart_metadata(metadata: &mut HashMap<String, S
}
fn should_persist_encryption_original_size(metadata: &HashMap<String, String>) -> bool {
metadata.keys().any(|key| is_object_encryption_marker(key))
metadata.keys().any(|key| is_encryption_metadata_key(key))
|| metadata.contains_key(SSEC_ALGORITHM_HEADER)
|| metadata.contains_key(SSEC_KEY_HEADER)
|| metadata.contains_key(SSEC_KEY_MD5_HEADER)
}
/// Per-set memoized capacity dirty scope.
@@ -686,6 +687,8 @@ mod core;
mod ctx;
mod metadata;
mod ops;
#[cfg(test)]
pub(crate) use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
#[cfg(feature = "test-util")]
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
pub(crate) use ops::object::body_cache_plaintext_len;
+23
View File
@@ -29,6 +29,12 @@ impl SetDisks {
pub async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
ListOperations::new(self.ctx()).delete_all(bucket, prefix).await
}
pub(crate) async fn delete_all_with_quorum(&self, bucket: &str, prefix: &str, write_quorum: usize) -> Result<()> {
ListOperations::new(self.ctx())
.delete_all_with_quorum(bucket, prefix, write_quorum)
.await
}
}
/// List/prefix maintenance operations, borrowing the `SetDisks` core state
@@ -48,6 +54,14 @@ impl<'a> ListOperations<'a> {
}
pub(crate) async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
self.delete_all_inner(bucket, prefix, None).await
}
async fn delete_all_with_quorum(&self, bucket: &str, prefix: &str, write_quorum: usize) -> Result<()> {
self.delete_all_inner(bucket, prefix, Some(write_quorum)).await
}
async fn delete_all_inner(&self, bucket: &str, prefix: &str, write_quorum: Option<usize>) -> Result<()> {
let disks = self.ctx.disks().read().await;
let disks = disks.clone();
@@ -79,6 +93,9 @@ impl<'a> ListOperations<'a> {
Ok(_) => {
errors.push(None);
}
Err(DiskError::FileNotFound | DiskError::PathNotFound | DiskError::VolumeNotFound) => {
errors.push(None);
}
Err(e) => {
errors.push(Some(e));
}
@@ -97,6 +114,12 @@ impl<'a> ListOperations<'a> {
);
}
if let Some(write_quorum) = write_quorum
&& let Some(err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum)
{
return Err(err.into());
}
Ok(())
}
}
+312 -15
View File
@@ -26,6 +26,8 @@ use crate::crash_inject::{self, CrashPoint};
use crate::multipart_listing::paginate_multipart_listing;
use futures::{StreamExt, stream};
use std::future::Future;
#[cfg(test)]
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use tokio::task::JoinSet;
@@ -33,7 +35,7 @@ const MULTIPART_LIST_IO_CONCURRENCY: usize = 16;
#[cfg(test)]
#[derive(Clone, Copy, PartialEq, Eq)]
enum MultipartCommitPause {
pub(crate) enum MultipartCommitPause {
PutPartBeforeLockLost,
PutPartAfterRename,
BeforeLockLost,
@@ -45,12 +47,13 @@ struct MultipartCommitBarrierState {
bucket: String,
object: String,
pause: MultipartCommitPause,
armed: AtomicBool,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
struct MultipartCommitBarrier {
pub(crate) struct MultipartCommitBarrier {
state: Arc<MultipartCommitBarrierState>,
}
@@ -60,11 +63,12 @@ static MULTIPART_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc
#[cfg(test)]
impl MultipartCommitBarrier {
fn install(bucket: &str, object: &str, pause: MultipartCommitPause) -> Self {
pub(crate) fn install(bucket: &str, object: &str, pause: MultipartCommitPause) -> Self {
let state = Arc::new(MultipartCommitBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
pause,
armed: AtomicBool::new(true),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
@@ -78,13 +82,13 @@ impl MultipartCommitBarrier {
Self { state }
}
async fn wait_until_paused(&self) {
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("multipart completion should reach the deterministic commit barrier");
}
fn release(&self) {
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
@@ -112,7 +116,9 @@ async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartComm
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object && barrier.pause == pause)
.cloned();
if let Some(barrier) = barrier {
if let Some(barrier) = barrier
&& barrier.armed.swap(false, Ordering::AcqRel)
{
barrier.arrived.notify_one();
barrier.release.notified().await;
}
@@ -777,6 +783,9 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
mut max_parts: usize,
opts: &ObjectOptions,
) -> Result<ListPartsInfo> {
let _upload_guard = self
.acquire_multipart_upload_read_lock("list_object_parts", bucket, object, upload_id, opts)
.await?;
let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, false).await?;
let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id);
@@ -1254,10 +1263,15 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
let _upload_guard = self
.acquire_multipart_upload_write_lock("abort_multipart_upload", bucket, object, upload_id, opts)
.await?;
self.check_upload_id_exists(bucket, object, upload_id, false).await?;
let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, true).await?;
let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id);
self.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await
self.delete_all_with_quorum(
RUSTFS_META_MULTIPART_BUCKET,
&upload_id_path,
fi.write_quorum(self.default_write_quorum()),
)
.await
}
// complete_multipart_upload finished
#[tracing::instrument(skip(self))]
@@ -1815,7 +1829,36 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
#[cfg(test)]
pause_multipart_commit(bucket, object, MultipartCommitPause::AfterRename).await;
drop(upload_guard);
let cleanup_store = self.clone();
let cleanup_upload_id_path = upload_id_path.clone();
let cleanup_bucket = bucket.to_owned();
let cleanup_object = object.to_owned();
let cleanup_upload_id = upload_id.to_owned();
let cleanup_handle = tokio::spawn(async move {
let _upload_guard = upload_guard;
if let Err(err) = cleanup_store
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &cleanup_upload_id_path, write_quorum)
.await
{
warn!(
bucket = %cleanup_bucket,
object = %cleanup_object,
upload_id = %cleanup_upload_id,
error = ?err,
"completed multipart upload staging cleanup did not reach write quorum"
);
}
});
if let Err(err) = cleanup_handle.await {
warn!(
bucket = %bucket,
object = %object,
upload_id = %upload_id,
error = ?err,
"completed multipart upload staging cleanup task failed"
);
}
drop(object_lock_guard); // drop object lock guard to release the lock
// backlog#1321: enqueue heal only when the committed replicas actually
@@ -1860,12 +1903,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
});
}
let upload_id_path = upload_id_path.clone();
let store = self.clone();
let _cleanup_handle = tokio::spawn(async move {
let _ = store.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await;
});
for (i, op_disk) in online_disks.iter().enumerate() {
if let Some(disk) = op_disk
&& disk.is_online().await
@@ -2177,6 +2214,122 @@ mod tests {
)
}
async fn assert_complete_first_linearizes(bucket: &'static str, object: &'static str, create_opts: ObjectOptions) {
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
make_bucket_on_all(&disk_stores, bucket).await;
let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x47; 4096], &create_opts).await;
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path));
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::AfterRename);
let complete_store = set_disks.clone();
let complete_upload_id = upload_id.clone();
let complete = tokio::spawn(async move {
complete_store
.complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default())
.await
});
barrier.wait_until_paused().await;
let abort_store = set_disks.clone();
let abort_upload_id = upload_id.clone();
let abort = tokio::spawn(async move {
abort_store
.abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default())
.await
});
signaling.wait_for_attempts(2).await;
assert!(!abort.is_finished(), "abort must wait for the completion upload lock");
barrier.release();
complete
.await
.expect("completion task should not panic")
.expect("completion should win the upload finalization");
let abort_err = abort
.await
.expect("abort task should not panic")
.expect_err("abort must observe the upload as finalized");
assert!(matches!(abort_err, StorageError::InvalidUploadID(..)));
set_disks
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("complete-first must leave the committed object readable");
assert!(matches!(
set_disks.check_upload_id_exists(bucket, object, &upload_id, false).await,
Err(StorageError::InvalidUploadID(..))
));
}
async fn assert_abort_first_linearizes(bucket: &'static str, object: &'static str, create_opts: ObjectOptions) {
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
make_bucket_on_all(&disk_stores, bucket).await;
let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x48; 4096], &create_opts).await;
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path.clone()));
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let object_holder = set_disks
.new_ns_lock(bucket, object)
.await
.expect("object namespace lock should be created")
.get_write_lock(Duration::from_secs(5))
.await
.expect("test should hold the object lock");
let holder = set_disks
.new_ns_lock(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path)
.await
.expect("upload namespace lock should be created")
.get_write_lock(Duration::from_secs(5))
.await
.expect("test should hold the upload lock");
signaling.wait_for_attempts(1).await;
let abort_store = set_disks.clone();
let abort_upload_id = upload_id.clone();
let abort = tokio::spawn(async move {
abort_store
.abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default())
.await
});
signaling.wait_for_attempts(2).await;
let complete_store = set_disks.clone();
let complete_upload_id = upload_id.clone();
let complete = tokio::spawn(async move {
complete_store
.complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default())
.await
});
drop(holder);
abort
.await
.expect("abort task should not panic")
.expect("abort should win the upload finalization");
drop(object_holder);
let complete_err = complete
.await
.expect("completion task should not panic")
.expect_err("completion must observe the aborted upload");
assert!(matches!(complete_err, StorageError::InvalidUploadID(..)));
let object_err = set_disks
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect_err("abort-first must not publish an object");
assert!(matches!(object_err, StorageError::ObjectNotFound(..)));
assert!(matches!(
set_disks.check_upload_id_exists(bucket, object, &upload_id, false).await,
Err(StorageError::InvalidUploadID(..))
));
}
async fn assert_quorum_minus_one_retry_preserves_completable_part(
disk_count: usize,
parity: usize,
@@ -3039,6 +3192,81 @@ mod tests {
.await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn abort_and_complete_linearize_for_plain_sse_and_legacy_layouts() {
assert_complete_first_linearizes("multipart-complete-first-plain", "object", ObjectOptions::default()).await;
assert_abort_first_linearizes("multipart-abort-first-plain", "object", ObjectOptions::default()).await;
let encrypted_opts = ObjectOptions {
user_defined: HashMap::from([(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string())]),
..Default::default()
};
temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async {
assert_complete_first_linearizes("multipart-complete-first-sse", "object", encrypted_opts.clone()).await;
assert_abort_first_linearizes("multipart-abort-first-sse", "object", encrypted_opts.clone()).await;
})
.await;
temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("false"))], async {
assert_complete_first_linearizes("multipart-complete-first-legacy", "object", encrypted_opts.clone()).await;
assert_abort_first_linearizes("multipart-abort-first-legacy", "object", encrypted_opts).await;
})
.await;
}
#[tokio::test]
async fn abort_enforces_delete_write_quorum_boundary() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "multipart-abort-delete-quorum";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
let quorum_upload = set_disks
.new_multipart_upload(bucket, object, &ObjectOptions::default())
.await
.expect("multipart upload should be created");
let saved_disks = {
let mut disks = set_disks.disks.write().await;
let saved = disks.clone();
disks[3] = None;
saved
};
set_disks
.abort_multipart_upload(bucket, object, &quorum_upload.upload_id, &ObjectOptions::default())
.await
.expect("abort should succeed at the exact delete write quorum");
*set_disks.disks.write().await = saved_disks;
assert!(matches!(
set_disks
.check_upload_id_exists(bucket, object, &quorum_upload.upload_id, false)
.await,
Err(StorageError::InvalidUploadID(..))
));
let below_quorum_upload = set_disks
.new_multipart_upload(bucket, object, &ObjectOptions::default())
.await
.expect("second multipart upload should be created");
let saved_disks = {
let mut disks = set_disks.disks.write().await;
let saved = disks.clone();
disks[2] = None;
disks[3] = None;
saved
};
let err = set_disks
.abort_multipart_upload(bucket, object, &below_quorum_upload.upload_id, &ObjectOptions::default())
.await
.expect_err("abort must report a delete below write quorum");
assert!(matches!(err, StorageError::ErasureWriteQuorum));
*set_disks.disks.write().await = saved_disks;
set_disks
.check_upload_id_exists(bucket, object, &below_quorum_upload.upload_id, false)
.await
.expect("failed abort must leave quorum-visible staging on the restored disks");
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn complete_revalidates_layout_candidate_after_upload_lock() {
@@ -3170,6 +3398,17 @@ mod tests {
tokio::task::yield_now().await;
assert!(!abort.is_finished(), "abort must wait until completion releases the upload lock");
let list_store = set_disks.clone();
let list_upload_id = upload_id.clone();
let list = tokio::spawn(async move {
list_store
.list_object_parts(bucket, object, &list_upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default())
.await
});
signaling.wait_for_attempts(3).await;
tokio::task::yield_now().await;
assert!(!list.is_finished(), "ListParts must wait until completion releases the upload lock");
barrier.release();
complete
.await
@@ -3180,10 +3419,68 @@ mod tests {
.expect("abort task should not panic")
.expect_err("the committed upload should no longer exist when abort acquires the lock");
assert!(matches!(abort_err, StorageError::InvalidUploadID(..)));
let list_err = list
.await
.expect("ListParts task should not panic")
.expect_err("the committed upload should no longer exist when ListParts acquires the lock");
assert!(matches!(list_err, StorageError::InvalidUploadID(..)));
})
.await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn complete_validates_parts_after_an_inflight_upload_part_commit() {
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
let bucket = "multipart-complete-put-part-race-bucket";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
let (upload_id, original_parts) =
stage_upload_with_create_opts(&set_disks, bucket, object, &[0x49; 4096], &ObjectOptions::default()).await;
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path));
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::PutPartBeforeLockLost);
let put_store = set_disks.clone();
let put_upload_id = upload_id.clone();
let put = tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(vec![0x4a; 4096]);
put_store
.put_object_part(bucket, object, &put_upload_id, 1, &mut reader, &ObjectOptions::default())
.await
});
barrier.wait_until_paused().await;
let complete_store = set_disks.clone();
let complete_upload_id = upload_id.clone();
let complete = tokio::spawn(async move {
complete_store
.complete_multipart_upload(bucket, object, &complete_upload_id, original_parts, &ObjectOptions::default())
.await
});
signaling.wait_for_attempts(2).await;
tokio::task::yield_now().await;
assert!(!complete.is_finished(), "completion must wait for the UploadPart commit lock");
barrier.release();
put.await
.expect("UploadPart task should not panic")
.expect("UploadPart replacement should commit");
let err = complete
.await
.expect("completion task should not panic")
.expect_err("completion must reject the stale ETag after UploadPart wins");
assert!(matches!(err, StorageError::InvalidPart(..)));
set_disks
.list_object_parts(bucket, object, &upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default())
.await
.expect("failed completion must leave the upload retryable");
}
#[tokio::test(start_paused = true)]
#[serial]
async fn complete_fences_upload_lock_loss_before_commit() {
+5 -71
View File
@@ -39,7 +39,6 @@ use crate::object_api::{GetObjectBodySource, get_object_body_cache_hook_suppress
use crate::services::tier::tier::{TierConfigMgr, TierOperationLease};
use crate::store::ECStore;
use futures::FutureExt as _;
use http::HeaderValue;
use std::future::Future;
fn erasure_from_file_info(fi: &FileInfo, uses_legacy: bool) -> Result<coding::Erasure> {
@@ -47,17 +46,6 @@ fn erasure_from_file_info(fi: &FileInfo, uses_legacy: bool) -> Result<coding::Er
.map_err(Error::from)
}
async fn get_object_reader_with_context(
ctx: &InstanceContext,
reader: Box<dyn AsyncRead + Unpin + Send + Sync>,
range: Option<HTTPRangeSpec>,
object_info: &ObjectInfo,
opts: &ObjectOptions,
headers: &HeaderMap<HeaderValue>,
) -> Result<(GetObjectReader, usize, i64)> {
GetObjectReader::new_with_resolver(reader, range, object_info, opts, headers, ctx.object_encryption_resolver()).await
}
/// Length of the full plaintext body when — and only when — this read's output
/// is exactly the object's complete plaintext, so the app-layer body cache may
/// serve it in place of the erasure read.
@@ -716,8 +704,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
size_bucket,
);
record_get_object_reader_path_observation(GET_OBJECT_PATH_CODEC_STREAMING, object_class, size_bucket);
let (mut reader, _offset, _length) =
get_object_reader_with_context(&self.ctx, stream, range, &object_info, opts, &h).await?;
let (mut reader, _offset, _length) = GetObjectReader::new(stream, range, &object_info, opts, &h).await?;
// Carry the hook probe result so the app layer skips its
// now-redundant lookup on the streaming miss path (ODC-16).
reader.body_source = body_source;
@@ -755,8 +742,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
let (rd, wd) = tokio::io::duplex(duplex_buffer_size);
debug!(bucket, object, duplex_buffer_size, "Created duplex pipe for object data transfer");
let (mut reader, offset, length) =
get_object_reader_with_context(&self.ctx, Box::new(rd), range, &object_info, opts, &h).await?;
let (mut reader, offset, length) = GetObjectReader::new(Box::new(rd), range, &object_info, opts, &h).await?;
// Carry the hook probe result so the app layer skips its now-redundant
// lookup on the streaming miss path (ODC-16).
reader.body_source = body_source;
@@ -4289,61 +4275,6 @@ mod erasure_construction_tests {
}
}
#[cfg(test)]
mod object_encryption_resolver_wiring_tests {
use super::*;
use crate::object_api::{EncryptionResolutionError, ObjectEncryptionResolver, ReadEncryptionMaterial, ReadEncryptionRequest};
use std::io::Cursor;
use std::sync::atomic::{AtomicUsize, Ordering};
struct CountingResolver {
calls: AtomicUsize,
}
#[async_trait::async_trait]
impl ObjectEncryptionResolver for CountingResolver {
async fn resolve_read_material(
&self,
_request: ReadEncryptionRequest<'_>,
) -> std::result::Result<Option<ReadEncryptionMaterial>, EncryptionResolutionError> {
self.calls.fetch_add(1, Ordering::Relaxed);
Ok(None)
}
}
#[tokio::test]
async fn get_object_reader_forwards_instance_resolver() {
let resolver = Arc::new(CountingResolver {
calls: AtomicUsize::new(0),
});
let ctx = InstanceContext::new();
assert!(
ctx.set_object_encryption_resolver(resolver.clone()).is_ok(),
"fresh context should accept resolver"
);
let object_info = ObjectInfo {
bucket: "bucket".to_string(),
name: "object".to_string(),
size: 1,
user_defined: Arc::new(HashMap::from([("x-amz-server-side-encryption".to_string(), "AES256".to_string())])),
..Default::default()
};
let result = get_object_reader_with_context(
&ctx,
Box::new(Cursor::new(Vec::<u8>::new())),
None,
&object_info,
&ObjectOptions::default(),
&HeaderMap::new(),
)
.await;
assert!(result.is_err(), "resolver returning no material must fail closed");
assert_eq!(resolver.calls.load(Ordering::Relaxed), 1);
}
}
#[cfg(test)]
pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
//! Shared hermetic `SetDisks` construction for the ops tests below: the
@@ -7093,6 +7024,9 @@ mod transition_source_identity_matrix_tests {
changed.metadata.insert("etag".to_string(), format!("changed-{index}"));
}
}
// Replace the object version list so VersionId drift removes the
// accepted source version instead of appending a second version.
changed.fresh = true;
for disk in &disk_stores {
disk.write_metadata("", bucket, &object, changed.clone())
.await
+17 -13
View File
@@ -1971,6 +1971,17 @@ mod metadata_cache_tests {
let mut fi = valid_test_fileinfo(object);
fi.mod_time = Some(OffsetDateTime::now_utc());
fi.erasure.index = fi.erasure.distribution[disk_index];
// `valid_test_fileinfo` carries a positive size with no parts —
// the shape `validate_collection_contents` rejects as
// `FileCorrupt` — so it can drive
// `get_object_with_fileinfo_rejects_positive_size_without_parts`.
// Metadata that has to survive a write needs the matching part.
fi.parts.push(ObjectPartInfo {
number: 1,
size: 1,
actual_size: 1,
..Default::default()
});
disk.write_metadata(bucket, bucket, object, fi)
.await
.expect("metadata should be written before quorum read");
@@ -2016,11 +2027,15 @@ mod metadata_cache_tests {
fi.size = 1;
fi.erasure.index = 1;
fi.metadata.insert("etag".to_string(), "etag-1".to_string());
fi.add_object_part(1, "part-etag".to_string(), 1, fi.mod_time, 1, None, None);
fi
}
#[tokio::test]
async fn get_object_with_fileinfo_rejects_positive_size_without_parts() {
let mut fi = valid_test_fileinfo("object");
fi.parts.clear();
let mut output = Vec::new();
let err = SetDisks::get_object_with_fileinfo(
"bucket",
@@ -2028,7 +2043,7 @@ mod metadata_cache_tests {
0,
1,
&mut output,
valid_test_fileinfo("object"),
fi,
Vec::new(),
&[],
0,
@@ -2119,12 +2134,6 @@ mod metadata_cache_tests {
let mut invalid_erasure = valid_test_fileinfo(object);
invalid_erasure.erasure.block_size = 0;
invalid_erasure.parts.push(ObjectPartInfo {
number: 1,
size: 1,
actual_size: 1,
..Default::default()
});
let err = SetDisks::get_object_with_fileinfo(
bucket,
object,
@@ -2159,6 +2168,7 @@ mod metadata_cache_tests {
let object = "empty";
let mut fi = valid_test_fileinfo(object);
fi.size = 0;
fi.parts.clear();
let mut output = Vec::new();
SetDisks::get_object_with_fileinfo(
@@ -2191,12 +2201,6 @@ mod metadata_cache_tests {
let mut fi = valid_test_fileinfo(object);
fi.erasure.block_size = 1;
fi.erasure.distribution = vec![1, 2, 3, 4];
fi.parts.push(ObjectPartInfo {
number: 1,
size: 1,
actual_size: 1,
..Default::default()
});
let mut output = Vec::new();
let err = SetDisks::get_object_with_fileinfo(
+4 -4
View File
@@ -2601,10 +2601,10 @@ mod tests {
// (backlog#1304): restore entry no longer serializes on the object lock.
// The replacement semantics — non-blocking reads during the copy-back and
// fast rejection of a concurrent restore — are covered end-to-end by
// `restore_object_usecase_reports_ongoing_conflict_and_completion`
// (rustfs/src/app/lifecycle_transition_api_test.rs) and at the lock level
// by the accept-guard test below; restore-vs-reader data protection lives
// in the inner put_object/complete_multipart_upload commit locks.
// `restore_object_usecase_reports_ongoing_conflict`
// (rustfs/src/app/lifecycle_transition_api_test.rs), while the SetDisks
// transition matrix covers the final local commit. Restore-vs-reader data
// protection lives in the inner put_object/complete_multipart_upload locks.
#[tokio::test]
#[serial_test::serial]
async fn restore_accept_guard_serializes_concurrent_accepts() {
+2 -2
View File
@@ -2,7 +2,7 @@
## MinIO-generated encrypted fixtures
`rustfs/src/storage/minio_generated_read_test.rs` validates the `bitrot -> GetObjectReader` path against raw MinIO backend data captured by
`minio_generated_read_test.rs` validates the `bitrot -> GetObjectReader` path against raw MinIO backend data captured by
`.\rustfs\scripts\minio_fixture_lab\lab.py`.
It currently covers multipart fixtures for:
@@ -20,5 +20,5 @@ Example:
```powershell
$env:RUSTFS_MINIO_FIXTURE_ROOT = '.\rustfs\tmp\minio-fixture-lab-local-key'
$env:RUSTFS_MINIO_STATIC_KMS_KEY_B64 = '<base64-32-byte-local-minio-kms-key>'
cargo +1.97.1 test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored
cargo +1.97.1 test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored
```
@@ -4,13 +4,14 @@ use std::fs;
use std::io::Cursor;
use std::path::{Path, PathBuf};
use super::sse::SseObjectEncryptionResolver;
use super::storage_api::ecstore_test_support::{
DiskAPI as _, DiskOption, Endpoint, Erasure, GetObjectReader, ObjectInfo, ObjectOptions, create_bitrot_reader, new_disk,
};
mod storage_api;
use rustfs_filemeta::{FileInfo, FileInfoOpts, get_file_info};
use serde::Deserialize;
use sha2::{Digest, Sha256};
use storage_api::minio_generated_read::{
DiskAPI as _, DiskOption, Endpoint, Erasure, GetObjectReader, ObjectInfo, ObjectOptions, create_bitrot_reader, new_disk,
};
use temp_env::async_with_vars;
use tokio::io::{AsyncReadExt, AsyncWrite};
@@ -48,7 +49,7 @@ impl AsyncWrite for VecAsyncWriter {
fn fixture_root() -> PathBuf {
std::env::var_os("RUSTFS_MINIO_FIXTURE_ROOT")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../crates/rio-v2/tests/fixtures/minio-generated"))
.unwrap_or_else(|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../rio-v2/tests/fixtures/minio-generated"))
}
fn case_dir(case_id: &str) -> PathBuf {
@@ -136,14 +137,12 @@ async fn read_fixture_plaintext(encrypted: Vec<u8>, object_info: ObjectInfo, kms
("RUSTFS_SSE_S3_MASTER_KEY", None::<String>),
],
async move {
let resolver = SseObjectEncryptionResolver;
let (mut reader, offset, length) = GetObjectReader::new_with_resolver(
let (mut reader, offset, length) = GetObjectReader::new(
Box::new(Cursor::new(encrypted)),
None,
&object_info,
&ObjectOptions::default(),
&http::HeaderMap::new(),
Some(&resolver),
)
.await
.map_err(|err| format!("construct GetObjectReader from MinIO raw fixture: {err:?}"))?;
@@ -220,12 +219,11 @@ async fn encrypted_fixture_bytes(case_dir: &Path, manifest: &ManifestRecord, fil
readers.push(reader);
}
let erasure = Erasure::try_new(
let erasure = Erasure::new(
file_info.erasure.data_blocks,
file_info.erasure.parity_blocks,
file_info.erasure.block_size,
)
.expect("fixture erasure geometry");
);
let mut writer = VecAsyncWriter::default();
let (written, err) = erasure.decode(&mut writer, readers, 0, part.size, part.size).await;
if let Some(err) = err {
+1 -1
View File
@@ -89,7 +89,7 @@ pub use error::{KmsError, KmsUnavailableError, Result};
pub use manager::KmsManager;
pub use service::{DataKey, ObjectEncryptionService};
pub use service_manager::{
KmsServiceManager, KmsServiceStatus, KmsStartOutcome, get_global_encryption_service, get_global_kms_service_manager,
KmsServiceManager, KmsServiceStatus, get_global_encryption_service, get_global_kms_service_manager,
init_global_kms_service_manager,
};
pub use types::*;
+89 -308
View File
@@ -20,12 +20,11 @@ use crate::error::{KmsError, Result};
use crate::manager::KmsManager;
use crate::service::ObjectEncryptionService;
use arc_swap::ArcSwap;
use std::future::Future;
use std::sync::{
Arc, OnceLock,
atomic::{AtomicU64, Ordering},
};
use tokio::sync::Mutex;
use tokio::sync::{Mutex, RwLock};
use tracing::{debug, error, info, warn};
const LOG_COMPONENT_KMS: &str = "kms";
@@ -45,13 +44,6 @@ pub enum KmsServiceStatus {
Error(String),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum KmsStartOutcome {
Started,
Restarted,
AlreadyRunning,
}
/// Service version information for zero-downtime reconfiguration
#[derive(Clone)]
struct ServiceVersion {
@@ -63,17 +55,16 @@ struct ServiceVersion {
manager: Arc<KmsManager>,
}
#[derive(Clone)]
struct RuntimeState {
config: Option<KmsConfig>,
status: KmsServiceStatus,
current_service: Option<ServiceVersion>,
}
/// Dynamic KMS service manager with versioned services for zero-downtime reconfiguration
pub struct KmsServiceManager {
/// Atomically published configuration, status, and current service.
state: ArcSwap<RuntimeState>,
/// Current service version (if running)
/// Uses ArcSwap for atomic, lock-free service switching
/// This allows instant atomic updates without blocking readers
current_service: ArcSwap<Option<ServiceVersion>>,
/// Current configuration
config: Arc<RwLock<Option<KmsConfig>>>,
/// Current status
status: Arc<RwLock<KmsServiceStatus>>,
/// Version counter (monotonically increasing)
version_counter: Arc<AtomicU64>,
/// Mutex to protect lifecycle operations (start, stop, reconfigure)
@@ -85,11 +76,9 @@ impl KmsServiceManager {
/// Create a new KMS service manager (not configured)
pub fn new() -> Self {
Self {
state: ArcSwap::from_pointee(RuntimeState {
config: None,
status: KmsServiceStatus::NotConfigured,
current_service: None,
}),
current_service: ArcSwap::from_pointee(None),
config: Arc::new(RwLock::new(None)),
status: Arc::new(RwLock::new(KmsServiceStatus::NotConfigured)),
version_counter: Arc::new(AtomicU64::new(0)),
lifecycle_mutex: Arc::new(Mutex::new(())),
}
@@ -97,65 +86,39 @@ impl KmsServiceManager {
/// Get current service status
pub async fn get_status(&self) -> KmsServiceStatus {
self.state.load().status.clone()
self.status.read().await.clone()
}
/// Get current configuration (if any)
pub async fn get_config(&self) -> Option<KmsConfig> {
self.state.load().config.clone()
self.config.read().await.clone()
}
/// Get configuration for status and management responses without static key material.
pub async fn get_redacted_config(&self) -> Option<KmsConfig> {
let mut config = self.state.load().config.clone()?;
Self::redact_config(&mut config);
Some(config)
}
/// Get status and redacted configuration from the same published snapshot.
pub async fn get_redacted_state(&self) -> (KmsServiceStatus, Option<KmsConfig>) {
let state = self.state.load();
let mut config = state.config.clone();
if let Some(config) = &mut config {
Self::redact_config(config);
}
(state.status.clone(), config)
}
fn redact_config(config: &mut KmsConfig) {
let mut config = self.config.read().await.clone()?;
if let BackendConfig::Static(static_config) = &mut config.backend_config {
use zeroize::Zeroize;
static_config.secret_key.zeroize();
}
Some(config)
}
/// Configure KMS with new configuration
pub async fn configure(&self, new_config: KmsConfig) -> Result<()> {
self.configure_with_persistence(new_config, || async { Ok(()) }).await
}
/// Configure KMS and publish the in-memory state only after persistence succeeds.
///
/// The persistence callback runs under the lifecycle lock and must not call
/// another lifecycle method on this manager.
pub async fn configure_with_persistence<Persist, PersistFuture>(&self, new_config: KmsConfig, persist: Persist) -> Result<()>
where
Persist: FnOnce() -> PersistFuture,
PersistFuture: Future<Output = Result<()>>,
{
new_config.validate()?;
let _guard = self.lifecycle_mutex.lock().await;
if self.state.load().current_service.is_some() {
return Err(KmsError::configuration_error(
"Cannot configure KMS while it is running; use reconfigure instead",
));
// Update configuration
{
let mut config = self.config.write().await;
*config = Some(new_config.clone());
}
// Update status
{
let mut status = self.status.write().await;
*status = KmsServiceStatus::Configured;
}
persist().await?;
self.state.store(Arc::new(RuntimeState {
config: Some(new_config),
status: KmsServiceStatus::Configured,
current_service: None,
}));
debug!(
event = EVENT_KMS_SERVICE_STATE,
@@ -173,35 +136,19 @@ impl KmsServiceManager {
self.start_internal().await
}
/// Start or restart KMS with the running-state decision serialized with the lifecycle action.
pub async fn start_or_restart(&self, force: bool) -> Result<KmsStartOutcome> {
let _guard = self.lifecycle_mutex.lock().await;
let running = self.state.load().current_service.is_some();
if running && !force {
return Ok(KmsStartOutcome::AlreadyRunning);
}
self.start_internal().await?;
Ok(if running {
KmsStartOutcome::Restarted
} else {
KmsStartOutcome::Started
})
}
/// Internal start implementation (called within lifecycle mutex)
async fn start_internal(&self) -> Result<()> {
let state = self.state.load_full();
let config = match state.config.as_ref() {
Some(config) => config.clone(),
None => {
let err_msg = "Cannot start KMS: no configuration provided";
error!("{}", err_msg);
self.state.store(Arc::new(RuntimeState {
config: None,
status: KmsServiceStatus::Error(err_msg.to_string()),
current_service: None,
}));
return Err(KmsError::configuration_error(err_msg));
let config = {
let config_guard = self.config.read().await;
match config_guard.as_ref() {
Some(config) => config.clone(),
None => {
let err_msg = "Cannot start KMS: no configuration provided";
error!("{}", err_msg);
let mut status = self.status.write().await;
*status = KmsServiceStatus::Error(err_msg.to_string());
return Err(KmsError::configuration_error(err_msg));
}
}
};
@@ -214,9 +161,17 @@ impl KmsServiceManager {
"KMS service starting"
);
match self.create_healthy_service_version(&config).await {
match self.create_service_version(&config).await {
Ok(service_version) => {
self.publish_running(config, service_version);
// Atomically update to new service version (lock-free, instant)
// ArcSwap::store() is a true atomic operation using CAS
self.current_service.store(Arc::new(Some(service_version)));
// Update status
{
let mut status = self.status.write().await;
*status = KmsServiceStatus::Running;
}
debug!(
event = EVENT_KMS_SERVICE_STATE,
@@ -230,24 +185,13 @@ impl KmsServiceManager {
Err(e) => {
let err_msg = format!("Failed to create KMS backend: {e}");
error!("{}", err_msg);
if state.current_service.is_none() {
self.state.store(Arc::new(RuntimeState {
config: state.config.clone(),
status: KmsServiceStatus::Error(err_msg.clone()),
current_service: None,
}));
}
let mut status = self.status.write().await;
*status = KmsServiceStatus::Error(err_msg.clone());
Err(KmsError::backend_error(&err_msg))
}
}
}
/// Replace the running service without exposing a stopped interval.
pub async fn restart(&self) -> Result<()> {
let _guard = self.lifecycle_mutex.lock().await;
self.start_internal().await
}
/// Stop KMS service
///
/// Note: This stops accepting new operations, but existing operations using
@@ -269,16 +213,15 @@ impl KmsServiceManager {
// Atomically clear current service version (lock-free, instant)
// Note: Existing Arc references will keep the service alive until operations complete
let state = self.state.load_full();
self.state.store(Arc::new(RuntimeState {
config: state.config.clone(),
status: if state.config.is_some() {
KmsServiceStatus::Configured
} else {
KmsServiceStatus::NotConfigured
},
current_service: None,
}));
self.current_service.store(Arc::new(None));
// Update status (keep configuration)
{
let mut status = self.status.write().await;
if !matches!(*status, KmsServiceStatus::NotConfigured) {
*status = KmsServiceStatus::Configured;
}
}
debug!(
event = EVENT_KMS_SERVICE_STATE,
@@ -301,22 +244,6 @@ impl KmsServiceManager {
/// This ensures zero downtime during reconfiguration, even for long-running
/// operations like encrypting large files.
pub async fn reconfigure(&self, new_config: KmsConfig) -> Result<()> {
self.reconfigure_with_persistence(new_config, || async { Ok(()) }).await
}
/// Reconfigure KMS after the candidate is healthy and persistence succeeds.
///
/// The persistence callback runs under the lifecycle lock and must not call
/// another lifecycle method on this manager.
pub async fn reconfigure_with_persistence<Persist, PersistFuture>(
&self,
new_config: KmsConfig,
persist: Persist,
) -> Result<()>
where
Persist: FnOnce() -> PersistFuture,
PersistFuture: Future<Output = Result<()>>,
{
let _guard = self.lifecycle_mutex.lock().await;
debug!(
@@ -328,16 +255,29 @@ impl KmsServiceManager {
);
new_config.validate()?;
// Configure with new config
{
let mut config = self.config.write().await;
*config = Some(new_config.clone());
}
// Create new service version without stopping old one
// This allows existing operations to continue while new operations use new service
match self.create_healthy_service_version(&new_config).await {
match self.create_service_version(&new_config).await {
Ok(new_service_version) => {
// Get old version for logging (lock-free read)
let old_version = self.state.load().current_service.as_ref().map(|sv| sv.version);
let old_version = self.current_service.load().as_ref().as_ref().map(|sv| sv.version);
persist().await?;
// Atomically switch to new service version (lock-free, instant CAS operation)
// This is a true atomic operation - no waiting for locks, instant switch
// Old service will be dropped when no more Arc references exist
self.current_service.store(Arc::new(Some(new_service_version.clone())));
self.publish_running(new_config, new_service_version.clone());
// Update status
{
let mut status = self.status.write().await;
*status = KmsServiceStatus::Running;
}
if let Some(old_ver) = old_version {
info!(
@@ -364,6 +304,8 @@ impl KmsServiceManager {
Err(e) => {
let err_msg = format!("Failed to reconfigure KMS: {e}");
error!("{}", err_msg);
let mut status = self.status.write().await;
*status = KmsServiceStatus::Error(err_msg.clone());
Err(KmsError::backend_error(&err_msg))
}
}
@@ -374,7 +316,7 @@ impl KmsServiceManager {
/// Returns the manager from the current service version.
/// Uses lock-free atomic load for optimal performance.
pub async fn get_manager(&self) -> Option<Arc<KmsManager>> {
self.state.load().current_service.as_ref().map(|sv| sv.manager.clone())
self.current_service.load().as_ref().as_ref().map(|sv| sv.manager.clone())
}
/// Get encryption service (if running)
@@ -384,7 +326,7 @@ impl KmsServiceManager {
/// This ensures new operations always use the latest service version,
/// while existing operations continue using their Arc references.
pub async fn get_encryption_service(&self) -> Option<Arc<ObjectEncryptionService>> {
self.state.load().current_service.as_ref().map(|sv| sv.service.clone())
self.current_service.load().as_ref().as_ref().map(|sv| sv.service.clone())
}
/// Get current service version number
@@ -392,16 +334,14 @@ impl KmsServiceManager {
/// Useful for monitoring and debugging.
/// Uses lock-free atomic load.
pub async fn get_service_version(&self) -> Option<u64> {
self.state.load().current_service.as_ref().map(|sv| sv.version)
self.current_service.load().as_ref().as_ref().map(|sv| sv.version)
}
/// Health check for the KMS service
pub async fn health_check(&self) -> Result<bool> {
let checked_state = self.state.load_full();
match checked_state.current_service.as_ref() {
Some(service_version) => {
let manager = service_version.manager.clone();
let checked_version = service_version.version;
let manager = self.get_manager().await;
match manager {
Some(manager) => {
// Perform health check on the backend
match manager.health_check().await {
Ok(healthy) => {
@@ -412,8 +352,9 @@ impl KmsServiceManager {
}
Err(e) => {
error!("KMS health check error: {}", e);
let _guard = self.lifecycle_mutex.lock().await;
self.mark_health_error_if_current(checked_version, &e);
// Update status to error
let mut status = self.status.write().await;
*status = KmsServiceStatus::Error(format!("Health check failed: {e}"));
Err(e)
}
}
@@ -472,33 +413,6 @@ impl KmsServiceManager {
manager: kms_manager,
})
}
async fn create_healthy_service_version(&self, config: &KmsConfig) -> Result<ServiceVersion> {
let service_version = self.create_service_version(config).await?;
if !service_version.manager.health_check().await? {
return Err(KmsError::backend_error("KMS backend health check failed"));
}
Ok(service_version)
}
fn publish_running(&self, config: KmsConfig, service_version: ServiceVersion) {
self.state.store(Arc::new(RuntimeState {
config: Some(config),
status: KmsServiceStatus::Running,
current_service: Some(service_version),
}));
}
fn mark_health_error_if_current(&self, checked_version: u64, error: &KmsError) {
let current = self.state.load_full();
if current.current_service.as_ref().map(|version| version.version) == Some(checked_version) {
self.state.store(Arc::new(RuntimeState {
config: current.config.clone(),
status: KmsServiceStatus::Error(format!("Health check failed: {error}")),
current_service: current.current_service.clone(),
}));
}
}
}
impl Default for KmsServiceManager {
@@ -531,11 +445,6 @@ pub async fn get_global_encryption_service() -> Option<Arc<ObjectEncryptionServi
#[cfg(test)]
mod tests {
use super::*;
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
fn static_config(key_id: &str, fill: u8) -> KmsConfig {
KmsConfig::static_kms(key_id.to_string(), BASE64_STANDARD.encode([fill; 32]))
}
#[tokio::test]
async fn configure_rejects_insecure_development_defaults_before_state_update() {
@@ -553,6 +462,8 @@ mod tests {
#[tokio::test]
async fn redacted_config_omits_static_key_material() {
use base64::Engine as _;
let manager = KmsServiceManager::new();
let encoded_key = base64::engine::general_purpose::STANDARD.encode([0x5au8; 32]);
manager
@@ -566,134 +477,4 @@ mod tests {
};
assert!(static_config.secret_key.is_empty());
}
#[tokio::test]
async fn configure_persistence_failure_leaves_state_unchanged() {
let manager = KmsServiceManager::new();
let result = manager
.configure_with_persistence(static_config("key-a", 0x11), || async { Err(KmsError::backend_error("persist failed")) })
.await;
assert!(result.is_err());
assert_eq!(manager.get_status().await, KmsServiceStatus::NotConfigured);
assert!(manager.get_config().await.is_none());
assert!(manager.get_encryption_service().await.is_none());
}
#[tokio::test]
async fn configure_rejects_running_service_without_changing_snapshot() {
let manager = KmsServiceManager::new();
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
manager.start().await.expect("start");
let version = manager.get_service_version().await;
let result = manager.configure(static_config("key-b", 0x22)).await;
assert!(result.is_err());
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
assert_eq!(manager.get_service_version().await, version);
assert_eq!(
manager.get_config().await.and_then(|config| config.default_key_id),
Some("key-a".to_string())
);
}
#[tokio::test]
async fn reconfigure_persistence_failure_keeps_old_running_snapshot() {
let manager = KmsServiceManager::new();
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
manager.start().await.expect("start");
let old_version = manager.get_service_version().await;
let old_service = manager.get_encryption_service().await.expect("old service");
let result = manager
.reconfigure_with_persistence(static_config("key-b", 0x22), || async {
Err(KmsError::backend_error("persist failed"))
})
.await;
assert!(result.is_err());
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
assert_eq!(manager.get_service_version().await, old_version);
assert_eq!(
manager.get_config().await.and_then(|config| config.default_key_id),
Some("key-a".to_string())
);
assert!(Arc::ptr_eq(
&old_service,
&manager.get_encryption_service().await.expect("old service remains")
));
}
#[tokio::test]
async fn reconfigure_candidate_failure_keeps_old_running_snapshot() {
let manager = KmsServiceManager::new();
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
manager.start().await.expect("start");
let old_version = manager.get_service_version().await;
let invalid_parent = tempfile::NamedTempFile::new().expect("temporary file");
let invalid_config = KmsConfig::local(invalid_parent.path().join("keys")).with_insecure_development_defaults();
let result = manager.reconfigure(invalid_config).await;
assert!(result.is_err());
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
assert_eq!(manager.get_service_version().await, old_version);
assert_eq!(
manager.get_config().await.and_then(|config| config.default_key_id),
Some("key-a".to_string())
);
}
#[tokio::test]
async fn restart_never_unpublishes_the_running_service() {
let manager = Arc::new(KmsServiceManager::new());
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
manager.start().await.expect("start");
let old_version = manager.get_service_version().await.expect("old version");
let restarting = {
let manager = manager.clone();
tokio::spawn(async move { manager.restart().await })
};
while !restarting.is_finished() {
assert!(manager.get_encryption_service().await.is_some());
tokio::task::yield_now().await;
}
restarting.await.expect("restart task").expect("restart");
assert!(manager.get_encryption_service().await.is_some());
assert!(manager.get_service_version().await.expect("new version") > old_version);
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
}
#[tokio::test]
async fn start_or_restart_decides_under_the_lifecycle_lock() {
let manager = KmsServiceManager::new();
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
assert_eq!(manager.start_or_restart(false).await.expect("initial start"), KmsStartOutcome::Started);
let first_version = manager.get_service_version().await.expect("first version");
assert_eq!(
manager.start_or_restart(false).await.expect("already running"),
KmsStartOutcome::AlreadyRunning
);
assert_eq!(manager.get_service_version().await, Some(first_version));
assert_eq!(manager.start_or_restart(true).await.expect("forced restart"), KmsStartOutcome::Restarted);
assert!(manager.get_service_version().await.expect("restarted version") > first_version);
}
#[tokio::test]
async fn stale_health_failure_cannot_poison_new_service_status() {
let manager = KmsServiceManager::new();
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
manager.start().await.expect("start");
let old_version = manager.get_service_version().await.expect("old version");
manager.restart().await.expect("restart");
manager.mark_health_error_if_current(old_version, &KmsError::backend_error("stale failure"));
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
}
}
+15
View File
@@ -25,6 +25,9 @@ keywords = ["file-system", "notification", "real-time", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-notify/latest/rustfs_notify/"
[features]
demo-examples = []
[dependencies]
rustfs-config = { workspace = true, features = ["notify", "server-config-model"] }
rustfs-ecstore = { workspace = true }
@@ -71,6 +74,18 @@ workspace = true
[lib]
doctest = false
[[example]]
name = "full_demo"
required-features = ["demo-examples"]
[[example]]
name = "full_demo_one"
required-features = ["demo-examples"]
[[example]]
name = "webhook"
required-features = ["demo-examples"]
[[bench]]
name = "snapshot_mode_scan"
harness = false
@@ -137,7 +137,7 @@ tests read):
./capture_via_docker.sh
RUSTFS_MINIO_STATIC_KMS_KEY_B64=IyqsU3kMFloCNup4BsZtf/rmfHVcTgznO2F25CkEH1g= \
cargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored
cargo test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored
```
This is exactly what the nightly `minio-interop` GitHub Actions workflow runs
+3 -73
View File
@@ -25,11 +25,6 @@ const RUSTFS_PREFIX: &str = "x-rustfs-";
const MINIO_PREFIX: &str = "x-minio-";
const MINIO_ENCRYPTION_PREFIX: &str = "x-minio-encryption-";
const RUSTFS_ENCRYPTION_PREFIX: &str = "x-rustfs-encryption-";
const MINIO_INTERNAL_ENCRYPTION_PREFIX: &str = "x-minio-internal-server-side-encryption-";
const MINIO_INTERNAL_ENCRYPTED_MULTIPART: &str = "x-minio-internal-encrypted-multipart";
const RUSTFS_ENCRYPTION_ORIGINAL_SIZE: &str = "x-rustfs-encryption-original-size";
const MINIO_ENCRYPTION_ORIGINAL_SIZE: &str = "x-minio-encryption-original-size";
const SSEC_ORIGINAL_SIZE: &str = "x-amz-server-side-encryption-customer-original-size";
// Suffix constants (part after x-rustfs- or x-minio-). Use with get_header/insert_header.
pub const SUFFIX_FORCE_DELETE: &str = "force-delete";
@@ -45,49 +40,11 @@ pub const SUFFIX_SOURCE_REPLICATION_REQUEST: &str = "source-replication-request"
pub const SUFFIX_SOURCE_REPLICATION_CHECK: &str = "source-replication-check";
pub const SUFFIX_REPLICATION_SSEC_CRC: &str = "replication-ssec-crc";
/// Returns true if the key is object-encryption metadata understood by RustFS or MinIO.
/// Case-insensitive for metadata filtering.
/// Returns true if the key is an internal encryption metadata key (x-rustfs-encryption-* or
/// x-minio-encryption-*). Case-insensitive for metadata filtering.
pub fn is_encryption_metadata_key(key: &str) -> bool {
let lower = key.to_lowercase();
lower.starts_with(RUSTFS_ENCRYPTION_PREFIX)
|| lower.starts_with(MINIO_ENCRYPTION_PREFIX)
|| lower.starts_with(MINIO_INTERNAL_ENCRYPTION_PREFIX)
|| lower == MINIO_INTERNAL_ENCRYPTED_MULTIPART
}
/// Returns true when a metadata key proves that object data is encrypted.
///
/// Original-size metadata alone is not proof: older plaintext objects can
/// retain that compatibility field after metadata migration.
pub fn is_object_encryption_marker(key: &str) -> bool {
(is_encryption_metadata_key(key)
&& !key.eq_ignore_ascii_case(RUSTFS_ENCRYPTION_ORIGINAL_SIZE)
&& !key.eq_ignore_ascii_case(MINIO_ENCRYPTION_ORIGINAL_SIZE))
|| super::is_sse_header(key)
}
/// Reads the logical object size recorded by encryption metadata.
pub fn get_object_encryption_original_size(metadata: &std::collections::HashMap<String, String>) -> std::io::Result<Option<i64>> {
let actual_size = super::get_str(metadata, super::SUFFIX_ACTUAL_SIZE);
let size = get_case_insensitive(metadata, RUSTFS_ENCRYPTION_ORIGINAL_SIZE)
.or_else(|| get_case_insensitive(metadata, SSEC_ORIGINAL_SIZE))
.or(actual_size.as_deref());
let Some(size) = size.filter(|size| !size.is_empty()) else {
return Ok(None);
};
size.parse::<i64>()
.map(Some)
.map_err(|error| std::io::Error::other(format!("Failed to parse encryption original size: {error}")))
}
fn get_case_insensitive<'a>(metadata: &'a std::collections::HashMap<String, String>, key: &str) -> Option<&'a str> {
metadata.get(key).map(String::as_str).or_else(|| {
metadata
.iter()
.find(|(candidate, _)| candidate.eq_ignore_ascii_case(key))
.map(|(_, value)| value.as_str())
})
lower.starts_with(RUSTFS_ENCRYPTION_PREFIX) || lower.starts_with(MINIO_ENCRYPTION_PREFIX)
}
fn rustfs_key(suffix: &str) -> String {
@@ -149,37 +106,10 @@ mod tests {
assert!(is_encryption_metadata_key("x-rustfs-encryption-iv"));
assert!(is_encryption_metadata_key("X-Rustfs-Encryption-Key"));
assert!(is_encryption_metadata_key("x-minio-encryption-iv"));
assert!(is_encryption_metadata_key("X-Minio-Internal-Server-Side-Encryption-Sealed-Key"));
assert!(is_encryption_metadata_key("X-Minio-Internal-Encrypted-Multipart"));
assert!(!is_encryption_metadata_key("x-amz-meta-custom"));
assert!(!is_encryption_metadata_key("x-rustfs-internal-healing"));
}
#[test]
fn object_encryption_marker_excludes_size_only_metadata() {
assert!(!is_object_encryption_marker(RUSTFS_ENCRYPTION_ORIGINAL_SIZE));
assert!(is_object_encryption_marker("X-Minio-Internal-Server-Side-Encryption-Sealed-Key"));
assert!(is_object_encryption_marker("x-amz-server-side-encryption"));
}
#[test]
fn object_encryption_original_size_is_case_insensitive() {
let metadata = std::collections::HashMap::from([(
"X-Amz-Server-Side-Encryption-Customer-Original-Size".to_string(),
"42".to_string(),
)]);
assert_eq!(get_object_encryption_original_size(&metadata).expect("valid size"), Some(42));
}
#[test]
fn object_encryption_original_size_prefers_rustfs_metadata() {
let metadata = std::collections::HashMap::from([
(SSEC_ORIGINAL_SIZE.to_string(), "21".to_string()),
(RUSTFS_ENCRYPTION_ORIGINAL_SIZE.to_string(), "42".to_string()),
]);
assert_eq!(get_object_encryption_original_size(&metadata).expect("valid size"), Some(42));
}
#[test]
fn test_get_header() {
let mut headers = HeaderMap::new();
@@ -358,7 +358,7 @@ Fixture-backed tests should run when the fixture path is present:
```bash
cargo test -p rustfs-ecstore --test legacy_bitrot_read_test -- --nocapture
cargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored --nocapture
cargo test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored --nocapture
```
## Multi-Expert Adversarial Review Summary
+179 -96
View File
@@ -363,27 +363,36 @@ impl Operation for ConfigureKmsHandler {
// Convert request to KmsConfig
let kms_config = configure_request.to_kms_config();
let persisted_config = kms_config.clone();
let (success, message, status) = match service_manager
.configure_with_persistence(kms_config, || async move {
save_kms_config(&persisted_config)
.await
.map_err(|error| rustfs_kms::KmsError::backend_error(format!("Failed to persist KMS configuration: {error}")))
})
.await
{
// Configure the service
let (success, message, status) = match service_manager.configure(kms_config.clone()).await {
Ok(()) => {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "configure",
state = "configured",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS configured successfully".to_string(), status)
// Persist the configuration to cluster storage
if let Err(e) = save_kms_config(&kms_config).await {
let error_msg = format!("KMS configured in memory but failed to persist: {e}");
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "configure",
state = "persist_failed",
error = %e,
"admin kms dynamic state"
);
let status = service_manager.get_status().await;
(false, error_msg, status)
} else {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "configure",
state = "configured",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS configured successfully".to_string(), status)
}
}
Err(e) => {
let error_msg = format!("Failed to configure KMS: {e}");
@@ -490,61 +499,125 @@ impl Operation for StartKmsHandler {
);
let service_manager = kms_service_manager_from_context();
let force = start_request.force.unwrap_or(false);
let (success, message, status) = match service_manager.start_or_restart(force).await {
Ok(rustfs_kms::KmsStartOutcome::Started) => {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "start",
state = "running",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS service started successfully".to_string(), status)
}
Ok(rustfs_kms::KmsStartOutcome::Restarted) => {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "restart",
state = "running",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS service restarted successfully".to_string(), status)
}
Ok(rustfs_kms::KmsStartOutcome::AlreadyRunning) => {
let status = service_manager.get_status().await;
warn!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "start",
state = "already_running",
"admin kms dynamic state"
);
(false, "KMS service is already running. Use force=true to restart.".to_string(), status)
}
Err(e) => {
let error_msg = format!("Failed to start or restart KMS service: {e}");
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "start",
state = "start_failed",
error = %e,
"admin kms dynamic state"
);
let status = service_manager.get_status().await;
(false, error_msg, status)
}
};
// Check if already running and force flag
let current_status = service_manager.get_status().await;
if matches!(current_status, KmsServiceStatus::Running) && !start_request.force.unwrap_or(false) {
warn!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "start",
state = "already_running",
"admin kms dynamic state"
);
let response = StartKmsResponse {
success: false,
message: "KMS service is already running. Use force=true to restart.".to_string(),
status: current_status,
};
let json_response = match serde_json::to_string(&response) {
Ok(json) => json,
Err(e) => {
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = EVENT_ADMIN_KMS_DYNAMIC_STATE,
operation = "start",
result = "response_serialize_failed",
error = %e,
"admin kms dynamic state"
);
return Ok(S3Response::new((
StatusCode::INTERNAL_SERVER_ERROR,
Body::from("Serialization error".to_string()),
)));
}
};
return Ok(S3Response::new((StatusCode::OK, Body::from(json_response))));
}
// Start the service (or restart if force=true)
let (success, message, status) =
if start_request.force.unwrap_or(false) && matches!(current_status, KmsServiceStatus::Running) {
// Force restart
match service_manager.stop().await {
Ok(()) => match service_manager.start().await {
Ok(()) => {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "restart",
state = "running",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS service restarted successfully".to_string(), status)
}
Err(e) => {
let error_msg = format!("Failed to restart KMS service: {e}");
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "restart",
state = "start_failed",
error = %e,
"admin kms dynamic state"
);
let status = service_manager.get_status().await;
(false, error_msg, status)
}
},
Err(e) => {
let error_msg = format!("Failed to stop KMS service for restart: {e}");
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "restart",
state = "stop_failed",
error = %e,
"admin kms dynamic state"
);
let status = service_manager.get_status().await;
(false, error_msg, status)
}
}
} else {
// Normal start
match service_manager.start().await {
Ok(()) => {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "start",
state = "running",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS service started successfully".to_string(), status)
}
Err(e) => {
let error_msg = format!("Failed to start KMS service: {e}");
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "start",
state = "start_failed",
error = %e,
"admin kms dynamic state"
);
let status = service_manager.get_status().await;
(false, error_msg, status)
}
}
};
let response = StartKmsResponse {
success,
@@ -701,7 +774,8 @@ impl Operation for GetKmsStatusHandler {
let service_manager = kms_service_manager_from_context();
let (status, config) = service_manager.get_redacted_state().await;
let status = service_manager.get_status().await;
let config = service_manager.get_redacted_config().await;
// Get backend type and health status
let backend_type = config.as_ref().map(|c| c.backend.clone());
@@ -834,27 +908,36 @@ impl Operation for ReconfigureKmsHandler {
// Convert request to KmsConfig
let kms_config = configure_request.to_kms_config();
let persisted_config = kms_config.clone();
let (success, message, status) = match service_manager
.reconfigure_with_persistence(kms_config, || async move {
save_kms_config(&persisted_config)
.await
.map_err(|error| rustfs_kms::KmsError::backend_error(format!("Failed to persist KMS configuration: {error}")))
})
.await
{
// Reconfigure the service (stops, reconfigures, and starts)
let (success, message, status) = match service_manager.reconfigure(kms_config.clone()).await {
Ok(()) => {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "reconfigure",
state = "reconfigured",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS reconfigured and restarted successfully".to_string(), status)
// Persist the configuration to cluster storage
if let Err(e) = save_kms_config(&kms_config).await {
let error_msg = format!("KMS reconfigured in memory but failed to persist: {e}");
error!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "reconfigure",
state = "persist_failed",
error = %e,
"admin kms dynamic state"
);
let status = service_manager.get_status().await;
(false, error_msg, status)
} else {
let status = service_manager.get_status().await;
info!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_KMS,
event = "kms_service_state",
operation = "reconfigure",
state = "reconfigured",
status = ?status,
"admin kms dynamic state"
);
(true, "KMS reconfigured and restarted successfully".to_string(), status)
}
}
Err(e) => {
let error_msg = format!("Failed to reconfigure KMS: {e}");
+8 -8
View File
@@ -18,12 +18,13 @@ use super::handles::{
IamHandle, KmsHandle, default_action_credential_interface, default_boot_time_interface, default_bucket_metadata_interface,
default_bucket_monitor_interface, default_buffer_config_interface, default_deployment_id_interface,
default_endpoints_interface, default_expiry_state_interface, default_federated_identity_interface,
default_internode_metrics_interface, default_local_node_name_interface, default_lock_client_interface,
default_lock_clients_interface, default_notification_system_interface, default_notify_interface,
default_outbound_tls_runtime_interface, default_performance_metrics_interface, default_region_interface,
default_replication_pool_interface, default_replication_stats_interface, default_runtime_port_interface,
default_s3select_db_interface, default_scanner_metrics_interface, default_server_config_interface,
default_storage_class_interface, default_tier_config_interface, default_transition_state_interface,
default_internode_metrics_interface, default_kms_runtime_interface, default_local_node_name_interface,
default_lock_client_interface, default_lock_clients_interface, default_notification_system_interface,
default_notify_interface, default_outbound_tls_runtime_interface, default_performance_metrics_interface,
default_region_interface, default_replication_pool_interface, default_replication_stats_interface,
default_runtime_port_interface, default_s3select_db_interface, default_scanner_metrics_interface,
default_server_config_interface, default_storage_class_interface, default_tier_config_interface,
default_transition_state_interface,
};
use super::interfaces::{
ActionCredentialInterface, BootTimeInterface, BucketMetadataInterface, BucketMonitorInterface, BufferConfigInterface,
@@ -79,7 +80,6 @@ pub struct AppContext {
impl AppContext {
pub fn new(object_store: Arc<ECStore>, iam: Arc<dyn IamInterface>, kms: Arc<dyn KmsInterface>) -> Self {
let object_data_cache = ObjectDataCacheAdapter::from_env_or_disabled();
let kms_runtime = Arc::new(crate::app::context::handles::KmsRuntimeHandle::new(kms.handle()));
// Let ecstore probe this cache inside get_object_reader, after
// metadata resolution but before the erasure data read (backlog#802).
crate::app::object_data_cache::register_object_data_cache_body_hook(Arc::clone(&object_data_cache));
@@ -94,7 +94,7 @@ impl AppContext {
iam,
federated_identity: default_federated_identity_interface(),
kms,
kms_runtime,
kms_runtime: default_kms_runtime_interface(),
outbound_tls_runtime: default_outbound_tls_runtime_interface(),
notify: default_notify_interface(),
notification_system: default_notification_system_interface(),
+7 -31
View File
@@ -128,19 +128,12 @@ impl KmsInterface for KmsHandle {
}
/// Default KMS runtime interface adapter.
pub struct KmsRuntimeHandle {
kms: Option<Arc<KmsServiceManager>>,
}
impl KmsRuntimeHandle {
pub fn new(kms: Arc<KmsServiceManager>) -> Self {
Self { kms: Some(kms) }
}
}
#[derive(Default)]
pub struct KmsRuntimeHandle;
impl KmsRuntimeInterface for KmsRuntimeHandle {
fn service_manager(&self) -> Option<Arc<KmsServiceManager>> {
self.kms.clone()
runtime_sources::kms_service_manager()
}
}
@@ -492,9 +485,7 @@ pub fn default_notification_system_interface() -> Arc<dyn NotificationSystemInte
}
pub fn default_kms_runtime_interface() -> Arc<dyn KmsRuntimeInterface> {
Arc::new(KmsRuntimeHandle {
kms: runtime_sources::kms_service_manager(),
})
Arc::new(KmsRuntimeHandle)
}
pub fn default_outbound_tls_runtime_interface() -> Arc<dyn OutboundTlsRuntimeInterface> {
@@ -616,10 +607,10 @@ pub fn default_buffer_config_interface() -> Arc<dyn BufferConfigInterface> {
#[cfg(test)]
mod tests {
use super::{
KmsRuntimeHandle, KmsServiceManager, ServerConfigHandle, default_federated_identity_interface,
federated_identity_interface, publish_default_federated_identity_service, runtime_sources,
ServerConfigHandle, default_federated_identity_interface, federated_identity_interface,
publish_default_federated_identity_service, runtime_sources,
};
use crate::app::context::interfaces::{KmsRuntimeInterface, ServerConfigInterface};
use crate::app::context::interfaces::ServerConfigInterface;
use rustfs_config::server_config::Config;
use rustfs_iam::{
federation::{FederatedIdentityRegistry, FederatedIdentityService, oidc::StandardOidcAdapter},
@@ -742,19 +733,4 @@ mod tests {
"handle B must serve its own credentials"
);
}
#[test]
fn kms_runtime_handles_keep_injected_managers_isolated() {
let manager_a = Arc::new(KmsServiceManager::new());
let manager_b = Arc::new(KmsServiceManager::new());
let handle_a = KmsRuntimeHandle::new(manager_a.clone());
let handle_b = KmsRuntimeHandle::new(manager_b.clone());
assert!(Arc::ptr_eq(&handle_a.service_manager().expect("manager A"), &manager_a));
assert!(Arc::ptr_eq(&handle_b.service_manager().expect("manager B"), &manager_b));
assert!(!Arc::ptr_eq(
&handle_a.service_manager().expect("manager A"),
&handle_b.service_manager().expect("manager B")
));
}
}
@@ -198,13 +198,9 @@ async fn data_usage_endpoint_serves_snapshot_without_live_listing() {
first_info.total_free_capacity
);
assert_eq!(
(first_info.total_capacity, first_info.total_free_capacity, first_info.total_used_capacity,),
(
second_info.total_capacity,
second_info.total_free_capacity,
second_info.total_used_capacity,
),
"repeated data usage requests must report stable capacity values"
second_info.total_used_capacity,
second_info.total_capacity.saturating_sub(second_info.total_free_capacity),
"server used capacity must stay internally consistent on repeated requests"
);
// The endpoint must serve the seeded snapshot numbers, not recomputed ones.
+80 -102
View File
@@ -169,6 +169,34 @@ async fn upload_test_object(ecstore: &Arc<ECStore>, bucket: &str, object: &str,
.expect("Failed to upload test object")
}
async fn transition_uploaded_object_directly(
ecstore: &Arc<ECStore>,
bucket: &str,
object: &str,
tier_name: &str,
uploaded: &ObjectInfo,
) -> ObjectInfo {
let transition_opts = ObjectOptions {
transition: lifecycle::lifecycle_contract::TransitionOptions {
status: lifecycle::lifecycle_contract::TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: uploaded.etag.clone().unwrap_or_default(),
..Default::default()
},
version_id: uploaded.version_id.map(|version| version.to_string()),
versioned: uploaded.version_id.is_some(),
mod_time: uploaded.mod_time,
..Default::default()
};
ecstore
.transition_object(bucket, object, &transition_opts)
.await
.expect("Failed to transition object directly");
wait_for_transition(ecstore, bucket, object, TRANSITION_WAIT_TIMEOUT)
.await
.expect("object should transition before restore assertions")
}
async fn set_bucket_lifecycle_transition_with_tier(
bucket_name: &str,
storage_class: &str,
@@ -2026,13 +2054,12 @@ async fn put_bucket_lifecycle_configuration_rejects_zero_day_expiration() {
/// backlog#1148 ilm-8: the RestoreObject API surface on a transitioned object.
///
/// POST restore(days=1) is accepted and immediately flips the object to
/// `x-amz-restore: ongoing-request="true"` (the mock tier's injected GET
/// latency keeps the background copy-back in flight); a second POST during
/// that window is rejected with 409 `RestoreAlreadyInProgress`; once the
/// copy-back completes the object reports `ongoing-request="false"` with a
/// future expiry-date; and a full GET is then served from the local restored
/// copy (the mock tier records no further `get` calls).
/// POST restore(days=1) is accepted and flips the object to
/// `x-amz-restore: ongoing-request="true"` while the mock tier GET barrier
/// proves the background copy-back has reached the remote read; a second POST
/// during that window is rejected with 409 `RestoreAlreadyInProgress`.
/// Synchronous SetDisks transition tests cover copy-back completion, restore
/// metadata, and local byte-identical reads.
///
/// Re-enabled in the serial lane by backlog#1304: the accept path now flips
/// the ongoing flag under a short compare-and-set guard and the copy-back
@@ -2042,7 +2069,7 @@ async fn put_bucket_lifecycle_configuration_rejects_zero_day_expiration() {
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial]
#[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#1148 (ilm-8)"]
async fn restore_object_usecase_reports_ongoing_conflict_and_completion() {
async fn restore_object_usecase_reports_ongoing_conflict() {
let (_disk_paths, ecstore) = setup_test_env().await;
let usecase = DefaultObjectUsecase::from_global();
@@ -2050,29 +2077,17 @@ async fn restore_object_usecase_reports_ongoing_conflict_and_completion() {
let backend = register_mock_tier(&tier_name).await;
let bucket = format!("test-api-restore-{}", &Uuid::new_v4().simple().to_string()[..8]);
// Must live under the `test/` prefix: `set_bucket_lifecycle_transition_with_tier`
// scopes the transition rule to `<Filter><Prefix>test/</Prefix>`, so an object
// outside it never matches, is never enqueued, and never transitions — the
// setup `wait_for_transition` would then time out before the restore assertions.
// Keep the object under the shared ILM test prefix even though this setup
// transitions it directly; it keeps diagnostics aligned with sibling tests.
let object = "test/restore/api-object.bin";
let payload: Vec<u8> = (0..128 * 1024).map(|i| (i % 251) as u8).collect();
create_test_bucket(&ecstore, bucket.as_str()).await;
set_bucket_lifecycle_transition_with_tier(bucket.as_str(), &tier_name)
.await
.expect("Failed to set lifecycle configuration");
let _ = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await;
let uploaded = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await;
let _ = transition_uploaded_object_directly(&ecstore, bucket.as_str(), object, &tier_name, &uploaded).await;
backend.clear_op_log().await;
lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(ecstore.clone(), bucket.as_str())
.await
.expect("Failed to enqueue transitioned object");
let _ = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT)
.await
.expect("object should transition before the restore API runs");
// Slow the tier GET so the background copy-back stays in flight long
// enough to observe the ongoing state and the conflict rejection.
backend.set_latency(Some(Duration::from_millis(1500))).await;
let get_barrier = backend.arm_get_barrier().await;
let restore_request = || RestoreRequest {
days: Some(1),
@@ -2096,15 +2111,18 @@ async fn restore_object_usecase_reports_ongoing_conflict_and_completion() {
.await
.expect("restore request should be accepted");
// The accepted restore is immediately visible as ongoing (the metadata is
// written synchronously before the copy-back is spawned).
get_barrier.wait_until_paused().await;
// The barrier proves the detached copy-back reached the tier GET and is
// still paused, so the ongoing state and conflict rejection are not timing
// assumptions about task scheduling.
let ongoing = ecstore
.get_object_info(bucket.as_str(), object, &ObjectOptions::default())
.await
.expect("Failed to load object info during restore");
assert!(
ongoing.restore_ongoing,
"x-amz-restore must report ongoing-request=true right after the restore is accepted"
"x-amz-restore must report ongoing-request=true while the copy-back tier GET is paused"
);
// A second restore while one is in flight is rejected.
@@ -2117,45 +2135,7 @@ async fn restore_object_usecase_reports_ongoing_conflict_and_completion() {
"unexpected rejection for a repeated restore: {err:?}"
);
// Completion: ongoing flips to false and a future expiry-date appears.
let mut completed = None;
for _ in 0..40 {
let info = ecstore
.get_object_info(bucket.as_str(), object, &ObjectOptions::default())
.await
.expect("Failed to poll object info for restore completion");
if !info.restore_ongoing && info.restore_expires.is_some() {
completed = Some(info);
break;
}
tokio::time::sleep(Duration::from_millis(500)).await;
}
let completed = completed.expect("restore copy-back should complete within the poll window");
backend.clear_faults().await;
let now_secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock before unix epoch")
.as_secs() as i64;
let expires = completed.restore_expires.expect("completed restore carries an expiry");
assert!(
expires.unix_timestamp() > now_secs,
"restore expiry-date must be in the future, got {expires}"
);
assert_eq!(
completed.transitioned_object.status, "complete",
"restore must not clear the transitioned state"
);
// The restored copy serves GET locally: no further tier GETs.
let tier_gets_after_restore = backend.get_count().await;
let data = read_object_bytes(&ecstore, bucket.as_str(), object).await;
assert_eq!(data, payload, "restored GET must return the original bytes");
assert_eq!(
backend.get_count().await,
tier_gets_after_restore,
"GET of a restored object must be served locally, not from the tier"
);
get_barrier.release();
}
/// backlog#1304: the restore-accept compare-and-set itself, under real
@@ -2176,27 +2156,18 @@ async fn restore_object_usecase_accepts_exactly_one_of_two_concurrent_restores()
let backend = register_mock_tier(&tier_name).await;
let bucket = format!("test-api-restore-cas-{}", &Uuid::new_v4().simple().to_string()[..8]);
// Must live under the `test/` prefix — the shared transition rule filters on it.
// Keep the object under the shared ILM test prefix for diagnostics parity.
let object = "test/restore/cas-object.bin";
let payload: Vec<u8> = (0..128 * 1024).map(|i| (i % 251) as u8).collect();
create_test_bucket(&ecstore, bucket.as_str()).await;
set_bucket_lifecycle_transition_with_tier(bucket.as_str(), &tier_name)
.await
.expect("Failed to set lifecycle configuration");
let _ = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await;
let uploaded = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await;
let _ = transition_uploaded_object_directly(&ecstore, bucket.as_str(), object, &tier_name, &uploaded).await;
backend.clear_op_log().await;
lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(ecstore.clone(), bucket.as_str())
.await
.expect("Failed to enqueue transitioned object");
let _ = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT)
.await
.expect("object should transition before the concurrent restores run");
// Keep the winner's copy-back in flight while the loser's accept runs, so
// the loser cannot slip into the already-restored path after a completed
// copy-back.
backend.set_latency(Some(Duration::from_millis(1500))).await;
// Hold the accepted copy-back at the tier GET until both accept attempts
// return, so the loser cannot observe an already-restored object.
let get_barrier = backend.arm_get_barrier().await;
let tier_gets_before_restore = backend.get_count().await;
@@ -2241,27 +2212,34 @@ async fn restore_object_usecase_accepts_exactly_one_of_two_concurrent_restores()
"the losing concurrent restore must be rejected as already in progress: {rejection:?}"
);
// Let the single accepted copy-back complete, then verify the tier saw
// exactly one restore read — a second GET means a double copy-back.
let mut completed = false;
for _ in 0..40 {
let info = ecstore
.get_object_info(bucket.as_str(), object, &ObjectOptions::default())
.await
.expect("Failed to poll object info for restore completion");
if !info.restore_ongoing && info.restore_expires.is_some() {
completed = true;
get_barrier.wait_until_paused().await;
get_barrier.release();
// This test is scoped to the accept CAS: completion, expiry metadata, and
// local restored GET service are covered by the single-request restore test
// above. Here it is enough to prove exactly one copy-back was admitted.
let expected_tier_gets = tier_gets_before_restore + 1;
let deadline = tokio::time::Instant::now() + TRANSITION_WAIT_TIMEOUT;
loop {
let actual_tier_gets = backend.get_count().await;
if actual_tier_gets >= expected_tier_gets {
assert_eq!(
actual_tier_gets - tier_gets_before_restore,
1,
"two concurrent restore requests must trigger exactly one tier copy-back GET"
);
break;
}
tokio::time::sleep(Duration::from_millis(500)).await;
if tokio::time::Instant::now() >= deadline {
let op_log = backend.op_log().await;
panic!(
"mock tier should record exactly one restore GET within {TRANSITION_WAIT_TIMEOUT:?}; \
tier_gets_before_restore={tier_gets_before_restore}, actual_tier_gets={actual_tier_gets}, op_log={op_log:?}"
);
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
backend.clear_faults().await;
assert!(completed, "the accepted restore copy-back should complete within the poll window");
assert_eq!(
backend.get_count().await - tier_gets_before_restore,
1,
"two concurrent restore requests must trigger exactly one tier copy-back GET"
);
}
/// rustfs/backlog#1320: a single PUT must compute the replication decision
+106 -18
View File
@@ -14,7 +14,7 @@
use super::{
module_switch::{resolve_notify_module_state, validate_notify_module_env, with_refreshed_notify_module_state_from},
refresh_persisted_module_switches_from_store, runtime_sources,
refresh_persisted_module_switches_from, runtime_sources,
};
use crate::storage_api::server::event::{
EventArgs as EcstoreEventArgs, StorageObjectInfo, read_existing_server_config_no_lock, register_event_dispatch_hook,
@@ -309,30 +309,36 @@ pub async fn shutdown_event_notifier() -> Result<(), NotificationError> {
#[instrument]
pub async fn init_event_notifier() -> Result<(), NotificationError> {
init_event_notifier_with_store(runtime_sources::current_object_store_handle).await
}
async fn init_event_notifier_with_store<CurrentStore>(current_store: CurrentStore) -> Result<(), NotificationError>
where
CurrentStore: FnOnce() -> Option<std::sync::Arc<rustfs_notify::NotifyStore>>,
{
mark_event_notifier_unreconciled();
validate_notify_module_env().map_err(NotificationError::Initialization)?;
let system = ensure_live_events_initialized();
refresh_persisted_module_switches_from_store()
let Some(store) = current_store() else {
let enabled = refresh_notify_module_enabled();
if enabled {
return Err(NotificationError::Initialization(
"failed to refresh notify module switch: storage layer not initialized".to_string(),
));
}
initialize_live_event_support(&system, None).await?;
return Ok(());
};
refresh_persisted_module_switches_from(store.clone())
.await
.map_err(|err| NotificationError::Initialization(format!("failed to refresh notify module switch: {err}")))?;
let enabled = refresh_notify_module_enabled();
if !enabled {
info!(
target: "rustfs::main::init_event_notifier",
"Notify module is disabled, initializing live event stream support only. Set {}=true to enable notification targets.",
rustfs_config::ENV_NOTIFY_ENABLE
);
if system.runtime_lifecycle_state() != NotificationRuntimeState::LiveOnly {
system.set_targets_enabled(false, None).await?;
}
system.reload_persisted_config().await?;
info!(
target: "rustfs::main::init_event_notifier",
"Live event stream support initialized successfully."
);
ensure_event_notifier_converged(&system)?;
initialize_live_event_support(&system, Some(store)).await?;
mark_event_notifier_reconciled();
return Ok(());
}
@@ -347,7 +353,7 @@ pub async fn init_event_notifier() -> Result<(), NotificationError> {
"Event notifier configuration found, proceeding with initialization."
);
system.reload_persisted_config().await?;
system.reload_persisted_config_from_store(store).await?;
let runtime_state = system.runtime_lifecycle_state();
if !matches!(runtime_state, NotificationRuntimeState::TargetsEnabled { .. }) {
system.set_targets_enabled(true, None).await?;
@@ -361,13 +367,53 @@ pub async fn init_event_notifier() -> Result<(), NotificationError> {
Ok(())
}
async fn initialize_live_event_support(
system: &NotificationSystem,
store: Option<std::sync::Arc<rustfs_notify::NotifyStore>>,
) -> Result<(), NotificationError> {
if store.is_some() {
info!(
target: "rustfs::main::init_event_notifier",
"Notify module is disabled, initializing live event stream support only. Set {}=true to enable notification targets.",
rustfs_config::ENV_NOTIFY_ENABLE
);
} else {
info!(
target: "rustfs::main::init_event_notifier",
"Notify module is disabled, initializing live event stream support only. Persisted notification target reconciliation is deferred until storage is initialized. Set {}=true to enable notification targets.",
rustfs_config::ENV_NOTIFY_ENABLE
);
}
if system.runtime_lifecycle_state() != NotificationRuntimeState::LiveOnly {
system.set_targets_enabled(false, None).await?;
}
if let Some(store) = store {
system.reload_persisted_config_from_store(store).await?;
}
info!(
target: "rustfs::main::init_event_notifier",
"Live event stream support initialized successfully."
);
ensure_event_notifier_converged(system)
}
#[cfg(test)]
mod tests {
use super::{convert_ecstore_object_info, parse_host_and_port, run_persisted_event_notifier_reconciler};
use super::super::module_switch::{
PersistedModuleSwitches, current_persisted_module_switches, persisted_module_switches_configured,
set_persisted_module_switches,
};
use super::{
convert_ecstore_object_info, init_event_notifier_with_store, parse_host_and_port, run_persisted_event_notifier_reconciler,
};
use crate::server::is_event_notifier_reconciled;
use crate::storage_api::server::event::StorageObjectInfo;
use crate::storage_api::server::event::contract::lifecycle::TransitionedObject;
use chrono::{DateTime, Utc};
use rustfs_notify::NotificationError;
use rustfs_notify::NotificationRuntimeState;
use serial_test::serial;
use std::{
collections::HashMap,
future::pending,
@@ -381,6 +427,48 @@ mod tests {
use tokio::sync::Notify;
use tokio_util::sync::CancellationToken;
#[tokio::test]
#[serial]
async fn disabled_notify_without_storage_initializes_live_events_without_reconcile() {
temp_env::async_with_vars([(rustfs_config::ENV_NOTIFY_ENABLE, Some("false"))], async {
let previous = current_persisted_module_switches();
let previous_configured = persisted_module_switches_configured();
set_persisted_module_switches(PersistedModuleSwitches::default(), false);
init_event_notifier_with_store(|| None)
.await
.expect("disabled notify should not require storage during early bootstrap");
let system = rustfs_notify::notification_system().expect("live event container should be initialized");
assert_eq!(system.runtime_lifecycle_state(), NotificationRuntimeState::LiveOnly);
assert!(
!is_event_notifier_reconciled(),
"early live-only bootstrap must not claim persisted notify config is reconciled"
);
set_persisted_module_switches(previous, previous_configured);
})
.await;
}
#[tokio::test]
#[serial]
async fn enabled_notify_without_storage_still_fails_visible() {
temp_env::async_with_vars([(rustfs_config::ENV_NOTIFY_ENABLE, Some("true"))], async {
let previous = current_persisted_module_switches();
let previous_configured = persisted_module_switches_configured();
set_persisted_module_switches(PersistedModuleSwitches::default(), false);
let err = init_event_notifier_with_store(|| None)
.await
.expect_err("enabled notify should still fail when storage is unavailable");
assert!(err.to_string().contains("storage layer not initialized"));
set_persisted_module_switches(previous, previous_configured);
})
.await;
}
#[test]
fn parse_host_and_port_with_ipv4_and_port() {
let (host, port) = parse_host_and_port("127.0.0.1:9000".to_string());
+1 -1
View File
@@ -80,7 +80,7 @@ pub(crate) fn current_persisted_module_switches() -> PersistedModuleSwitches {
}
}
fn persisted_module_switches_configured() -> bool {
pub(crate) fn persisted_module_switches_configured() -> bool {
PERSISTED_MODULE_SWITCH_CONFIGURED.load(Ordering::Relaxed)
}
-2
View File
@@ -36,8 +36,6 @@ mod ecfs_extend;
mod ecfs_test;
pub(crate) mod head_prefix;
#[cfg(test)]
mod minio_generated_read_test;
#[cfg(test)]
mod multi_factor_scheduler_integration_test;
pub(crate) mod runtime_sources;
#[cfg(test)]
+13 -270
View File
@@ -70,10 +70,6 @@
//! ```
use super::StorageError;
use super::storage_api::ecstore_object::{
EncryptionResolutionError, EncryptionResolutionErrorKind, ObjectEncryptionResolver, ReadEncryptionMaterial,
ReadEncryptionMode, ReadEncryptionRequest,
};
use crate::storage::storage_api::runtime_sources_consumer::runtime_sources;
#[cfg(feature = "rio-v2")]
use aes_gcm::aead::Payload;
@@ -148,7 +144,6 @@ use rustfs_utils::http::headers::{
};
use rustfs_utils::path::path_join_buf;
use s3s::dto::{SSECustomerAlgorithm, SSECustomerKey, SSECustomerKeyMD5, SSEKMSKeyId};
use std::borrow::Cow;
// ============================================================================
// High-Level SSE Configuration
@@ -646,23 +641,6 @@ pub(crate) fn validate_sse_headers_for_read(metadata: &HashMap<String, String>,
}
pub(crate) fn map_get_object_reader_error(err: StorageError) -> ApiError {
if let StorageError::Io(io_error) = &err
&& let Some(resolution_error) = io_error
.get_ref()
.and_then(|source| source.downcast_ref::<EncryptionResolutionError>())
{
let code = match resolution_error.kind() {
EncryptionResolutionErrorKind::InvalidRequest => S3ErrorCode::InvalidRequest,
EncryptionResolutionErrorKind::ServiceUnavailable => S3ErrorCode::ServiceUnavailable,
_ => S3ErrorCode::InternalError,
};
return ApiError {
code,
message: resolution_error.to_string(),
source: Some(Box::new(err)),
};
}
if let Some(message) = map_ssec_get_object_reader_error_message(&err) {
return ApiError {
code: S3ErrorCode::InvalidRequest,
@@ -785,104 +763,6 @@ pub enum EncryptionKeyKind {
Object,
}
pub(crate) struct SseObjectEncryptionResolver;
#[async_trait]
impl ObjectEncryptionResolver for SseObjectEncryptionResolver {
async fn resolve_read_material(
&self,
request: ReadEncryptionRequest<'_>,
) -> Result<Option<ReadEncryptionMaterial>, EncryptionResolutionError> {
let metadata = normalize_encryption_metadata_case(request.metadata)?;
let (_, customer_key, customer_key_md5) =
extract_ssec_params_from_headers(request.headers).map_err(map_encryption_resolution_error)?;
let material = sse_decryption(DecryptionRequest {
bucket: request.bucket,
key: request.object,
metadata: &metadata,
sse_customer_key: customer_key.as_ref(),
sse_customer_key_md5: customer_key_md5.as_ref(),
})
.await
.map_err(map_encryption_resolution_error)?;
Ok(material.map(|material| ReadEncryptionMaterial {
key_bytes: material.key_bytes,
mode: match material.key_kind {
EncryptionKeyKind::Direct => ReadEncryptionMode::Direct {
base_nonce: material.base_nonce,
},
EncryptionKeyKind::Object => ReadEncryptionMode::Object,
},
}))
}
}
fn normalize_encryption_metadata_case(
metadata: &HashMap<String, String>,
) -> Result<Cow<'_, HashMap<String, String>>, EncryptionResolutionError> {
const CANONICAL_KEYS: &[&str] = &[
"x-amz-server-side-encryption",
"x-amz-server-side-encryption-aws-kms-key-id",
"x-amz-server-side-encryption-customer-algorithm",
"x-amz-server-side-encryption-customer-key-md5",
SSEC_ORIGINAL_SIZE_HEADER,
INTERNAL_ENCRYPTION_KEY_ID_HEADER,
INTERNAL_ENCRYPTION_KEY_HEADER,
INTERNAL_ENCRYPTION_ALGORITHM_HEADER,
INTERNAL_ENCRYPTION_IV_HEADER,
"x-rustfs-encryption-context",
"x-rustfs-encryption-tag",
INTERNAL_ENCRYPTION_ORIGINAL_SIZE_HEADER,
MINIO_INTERNAL_ENCRYPTION_MULTIPART_HEADER,
MINIO_INTERNAL_ENCRYPTION_IV_HEADER,
MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER,
MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER,
MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER,
MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER,
MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER,
MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER,
];
let needs_normalization = metadata.keys().any(|key| {
CANONICAL_KEYS
.iter()
.any(|canonical| key != canonical && key.eq_ignore_ascii_case(canonical))
});
if !needs_normalization {
return Ok(Cow::Borrowed(metadata));
}
let mut normalized = metadata.clone();
for canonical in CANONICAL_KEYS {
let mut matching_values = metadata
.iter()
.filter_map(|(key, value)| key.eq_ignore_ascii_case(canonical).then_some(value));
let Some(value) = matching_values.next() else {
continue;
};
if matching_values.any(|candidate| candidate != value) {
return Err(EncryptionResolutionError::new(
EncryptionResolutionErrorKind::InvalidMetadata,
format!("conflicting object encryption metadata for {canonical}"),
));
}
if !normalized.contains_key(*canonical) {
normalized.insert((*canonical).to_string(), value.clone());
}
}
Ok(Cow::Owned(normalized))
}
fn map_encryption_resolution_error(error: ApiError) -> EncryptionResolutionError {
let kind = match error.code {
S3ErrorCode::InvalidArgument | S3ErrorCode::InvalidRequest => EncryptionResolutionErrorKind::InvalidRequest,
S3ErrorCode::ServiceUnavailable => EncryptionResolutionErrorKind::ServiceUnavailable,
_ => EncryptionResolutionErrorKind::DecryptionFailed,
};
EncryptionResolutionError::new(kind, error.message)
}
#[derive(Debug, Clone)]
pub struct ManagedSealedKey {
#[cfg(feature = "rio-v2")]
@@ -2751,20 +2631,19 @@ fn ssec_invalid_request(message: &str) -> ApiError {
mod tests {
use super::{
ApiError, DataKey, DecryptionRequest, EncryptionKeyKind, EncryptionMaterial, EncryptionRequest,
EncryptionResolutionErrorKind, INTERNAL_ENCRYPTION_ALGORITHM_HEADER, INTERNAL_ENCRYPTION_IV_HEADER,
INTERNAL_ENCRYPTION_KEY_HEADER, INTERNAL_ENCRYPTION_KEY_ID_HEADER, KmsSseDekProvider, KmsUnavailableError,
MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER, MINIO_INTERNAL_ENCRYPTION_IV_HEADER,
MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER, MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER,
MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER, MINIO_INTERNAL_ENCRYPTION_MULTIPART_HEADER,
MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER, MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER,
ObjectEncryptionResolver, PrepareEncryptionRequest, ReadEncryptionMode, ReadEncryptionRequest, SSEC_ORIGINAL_SIZE_HEADER,
SSEType, SseDekProvider, SseObjectEncryptionResolver, SsecParams, StorageError, TestSseDekProvider,
apply_managed_decryption_material, apply_managed_encryption_material, encryption_material_to_metadata,
extract_server_side_encryption_from_headers, extract_ssec_params_from_headers, extract_ssekms_context_from_headers,
generate_ssec_nonce, is_managed_sse, kms_operation_error, map_get_object_reader_error, mark_encrypted_multipart_metadata,
normalize_managed_metadata, reset_sse_dek_provider, resolve_effective_kms_key_id, sse_decryption, sse_encryption,
sse_prepare_encryption, strip_managed_encryption_metadata, validate_sse_headers_for_read, validate_sse_headers_for_write,
validate_ssec_for_read, validate_ssec_params, verify_ssec_key_match,
INTERNAL_ENCRYPTION_ALGORITHM_HEADER, INTERNAL_ENCRYPTION_IV_HEADER, INTERNAL_ENCRYPTION_KEY_HEADER,
INTERNAL_ENCRYPTION_KEY_ID_HEADER, KmsSseDekProvider, KmsUnavailableError, MINIO_INTERNAL_ENCRYPTION_ALGORITHM_HEADER,
MINIO_INTERNAL_ENCRYPTION_IV_HEADER, MINIO_INTERNAL_ENCRYPTION_KMS_CONTEXT_HEADER,
MINIO_INTERNAL_ENCRYPTION_KMS_KEY_ID_HEADER, MINIO_INTERNAL_ENCRYPTION_KMS_SEALED_KEY_HEADER,
MINIO_INTERNAL_ENCRYPTION_MULTIPART_HEADER, MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER,
MINIO_INTERNAL_ENCRYPTION_SSEC_SEALED_KEY_HEADER, PrepareEncryptionRequest, SSEC_ORIGINAL_SIZE_HEADER, SSEType,
SseDekProvider, SsecParams, StorageError, TestSseDekProvider, apply_managed_decryption_material,
apply_managed_encryption_material, encryption_material_to_metadata, extract_server_side_encryption_from_headers,
extract_ssec_params_from_headers, extract_ssekms_context_from_headers, generate_ssec_nonce, is_managed_sse,
kms_operation_error, map_get_object_reader_error, mark_encrypted_multipart_metadata, normalize_managed_metadata,
reset_sse_dek_provider, resolve_effective_kms_key_id, sse_decryption, sse_encryption, sse_prepare_encryption,
strip_managed_encryption_metadata, validate_sse_headers_for_read, validate_sse_headers_for_write, validate_ssec_for_read,
validate_ssec_params, verify_ssec_key_match,
};
#[cfg(feature = "rio-v2")]
use super::{
@@ -2824,97 +2703,6 @@ mod tests {
SSE_TEST_LOCK.get_or_init(|| Mutex::new(())).lock().await
}
#[tokio::test]
async fn object_encryption_resolver_returns_ssec_read_material() {
let key = [0x31; 32];
let key_b64 = BASE64_STANDARD.encode(key);
let key_md5 = BASE64_STANDARD.encode(md5::compute(key).0);
let nonce = [0x42; 12];
let metadata = HashMap::from([
("X-Amz-Server-Side-Encryption-Customer-Algorithm".to_string(), "AES256".to_string()),
("X-Amz-Server-Side-Encryption-Customer-Key-Md5".to_string(), key_md5.clone()),
("X-Rustfs-Encryption-Iv".to_string(), BASE64_STANDARD.encode(nonce)),
]);
let mut headers = HeaderMap::new();
headers.insert("x-amz-server-side-encryption-customer-algorithm", HeaderValue::from_static("AES256"));
headers.insert(
"x-amz-server-side-encryption-customer-key",
HeaderValue::from_str(&key_b64).expect("base64 key is a valid header"),
);
headers.insert(
"x-amz-server-side-encryption-customer-key-md5",
HeaderValue::from_str(&key_md5).expect("base64 MD5 is a valid header"),
);
let material = SseObjectEncryptionResolver
.resolve_read_material(ReadEncryptionRequest {
bucket: "bucket",
object: "object",
metadata: &metadata,
headers: &headers,
})
.await
.expect("SSE-C material should resolve")
.expect("SSE-C metadata should produce material");
assert_eq!(material.key_bytes, key);
assert_eq!(material.mode, ReadEncryptionMode::Direct { base_nonce: nonce });
}
#[tokio::test]
async fn object_encryption_resolver_classifies_missing_ssec_key_as_invalid_request() {
let metadata = HashMap::from([("x-amz-server-side-encryption-customer-algorithm".to_string(), "AES256".to_string())]);
let result = SseObjectEncryptionResolver
.resolve_read_material(ReadEncryptionRequest {
bucket: "bucket",
object: "object",
metadata: &metadata,
headers: &HeaderMap::new(),
})
.await;
let error = match result {
Err(error) => error,
Ok(_) => panic!("missing SSE-C key must fail closed"),
};
assert_eq!(error.kind(), EncryptionResolutionErrorKind::InvalidRequest);
}
#[tokio::test]
async fn object_encryption_resolver_rejects_conflicting_metadata_case_variants() {
let metadata = HashMap::from([
("x-rustfs-encryption-key".to_string(), "first".to_string()),
("X-Rustfs-Encryption-Key".to_string(), "second".to_string()),
]);
let result = SseObjectEncryptionResolver
.resolve_read_material(ReadEncryptionRequest {
bucket: "bucket",
object: "object",
metadata: &metadata,
headers: &HeaderMap::new(),
})
.await;
let error = match result {
Err(error) => error,
Ok(_) => panic!("conflicting metadata aliases must fail closed"),
};
assert_eq!(error.kind(), EncryptionResolutionErrorKind::InvalidMetadata);
}
#[test]
fn normalize_encryption_metadata_case_accepts_lowercase_minio_internal_keys() {
let lowercase_key = MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER.to_ascii_lowercase();
let metadata = HashMap::from([(lowercase_key, "sealed-key".to_string())]);
let normalized = super::normalize_encryption_metadata_case(&metadata).expect("metadata aliases should normalize");
assert_eq!(
normalized.get(MINIO_INTERNAL_ENCRYPTION_S3_SEALED_KEY_HEADER),
Some(&"sealed-key".to_string())
);
}
struct UnavailableSseDekProvider;
#[async_trait::async_trait]
@@ -3937,19 +3725,6 @@ mod tests {
assert_eq!(decrypted.key_kind, EncryptionKeyKind::Object);
assert_eq!(decrypted.key_bytes, material.key_bytes);
let resolved = SseObjectEncryptionResolver
.resolve_read_material(ReadEncryptionRequest {
bucket: "bucket",
object: "object",
metadata: &metadata,
headers: &HeaderMap::new(),
})
.await
.expect("managed resolver")
.expect("managed material");
assert_eq!(resolved.mode, ReadEncryptionMode::Object);
assert_eq!(resolved.key_bytes, material.key_bytes);
},
)
.await;
@@ -4010,29 +3785,6 @@ mod tests {
assert_eq!(decrypted.key_kind, EncryptionKeyKind::Object);
assert_eq!(decrypted.key_bytes, material.key_bytes);
let mut headers = HeaderMap::new();
headers.insert("x-amz-server-side-encryption-customer-algorithm", HeaderValue::from_static("AES256"));
headers.insert(
"x-amz-server-side-encryption-customer-key",
HeaderValue::from_str(&customer_key).expect("customer key header"),
);
headers.insert(
"x-amz-server-side-encryption-customer-key-md5",
HeaderValue::from_str(&customer_key_md5).expect("customer key MD5 header"),
);
let resolved = SseObjectEncryptionResolver
.resolve_read_material(ReadEncryptionRequest {
bucket: "bucket",
object: "object",
metadata: &metadata,
headers: &headers,
})
.await
.expect("SSE-C resolver")
.expect("SSE-C material");
assert_eq!(resolved.mode, ReadEncryptionMode::Object);
assert_eq!(resolved.key_bytes, material.key_bytes);
}
#[cfg(feature = "rio-v2")]
@@ -4963,15 +4715,6 @@ mod tests {
);
}
#[test]
fn test_map_get_object_reader_error_preserves_typed_service_unavailable() {
let resolution_error =
super::EncryptionResolutionError::new(EncryptionResolutionErrorKind::ServiceUnavailable, "KMS unavailable");
let err = map_get_object_reader_error(StorageError::other(resolution_error));
assert_eq!(err.code, S3ErrorCode::ServiceUnavailable);
assert_eq!(err.message, "KMS unavailable");
}
#[test]
fn test_map_get_object_reader_error_leaves_non_ssec_errors_unchanged() {
let err = map_get_object_reader_error(StorageError::other("plain io failure"));
+6 -33
View File
@@ -510,21 +510,12 @@ pub(crate) mod ecstore_object {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::object::GetObjectBodySource;
pub(crate) use rustfs_ecstore::api::object::{
EncryptionResolutionError, EncryptionResolutionErrorKind, GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup,
ObjectEncryptionResolver, ObjectMutationHook, ReadEncryptionMaterial, ReadEncryptionMode, ReadEncryptionRequest,
get_object_body_cache_plaintext_len, lookup_get_object_body_cache_hook, register_get_object_body_cache_hook,
register_object_mutation_hook, unregister_get_object_body_cache_hook, unregister_object_mutation_hook,
GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, ObjectMutationHook, get_object_body_cache_plaintext_len,
lookup_get_object_body_cache_hook, register_get_object_body_cache_hook, register_object_mutation_hook,
unregister_get_object_body_cache_hook, unregister_object_mutation_hook,
};
}
#[cfg(all(test, feature = "rio-v2"))]
pub(crate) mod ecstore_test_support {
pub(crate) use rustfs_ecstore::api::bitrot::create_bitrot_reader;
pub(crate) use rustfs_ecstore::api::disk::{DiskAPI, DiskOption, endpoint::Endpoint, new_disk};
pub(crate) use rustfs_ecstore::api::erasure::Erasure;
pub(crate) use rustfs_ecstore::api::object::{GetObjectReader, ObjectInfo, ObjectOptions};
}
pub(crate) mod ecstore_set_disk {
pub(crate) use rustfs_ecstore::api::set_disk::{DEFAULT_READ_BUFFER_SIZE, get_lock_acquire_timeout, is_valid_storage_class};
}
@@ -954,21 +945,13 @@ pub(crate) async fn init_local_disks(endpoint_pools: EndpointServerPools) -> Res
/// The process-level bootstrap instance context that single-instance startup
/// threads through the storage foundation (Phase 5 follow-up, backlog#1052).
pub(crate) fn bootstrap_instance_ctx() -> Arc<InstanceContext> {
let context = ecstore_runtime::bootstrap_ctx();
configure_object_encryption_resolver(&context);
context
ecstore_runtime::bootstrap_ctx()
}
/// Construct a fresh per-server instance context (backlog#1052 S5): a second
/// embedded server owns its own erasure/region/endpoint/deployment id cells.
pub(crate) fn new_instance_ctx() -> Arc<InstanceContext> {
let context = Arc::new(InstanceContext::new());
configure_object_encryption_resolver(&context);
context
}
fn configure_object_encryption_resolver(context: &InstanceContext) {
let _ = context.set_object_encryption_resolver(Arc::new(super::sse::SseObjectEncryptionResolver));
Arc::new(InstanceContext::new())
}
pub(crate) fn init_lock_clients(endpoint_pools: EndpointServerPools) {
@@ -1730,7 +1713,7 @@ pub(crate) async fn init_compression_total_memory_from_backend(store: Arc<ECStor
mod tests {
use super::{
apply_active_resync_intents, bucket_targets_metadata_lock_shard, ecstore_bucket, lock_bucket_targets_metadata,
new_instance_ctx, scanner_maintenance_config_file,
scanner_maintenance_config_file,
};
use std::time::Duration;
@@ -1761,16 +1744,6 @@ mod tests {
);
}
#[test]
fn fresh_instance_context_installs_object_encryption_resolver() {
assert!(new_instance_ctx().object_encryption_resolver().is_some());
}
#[test]
fn bootstrap_instance_context_installs_object_encryption_resolver() {
assert!(super::bootstrap_instance_ctx().object_encryption_resolver().is_some());
}
#[test]
fn scanner_maintenance_config_only_includes_scanner_owned_work() {
assert!(scanner_maintenance_config_file(ecstore_bucket::metadata::BUCKET_LIFECYCLE_CONFIG));
+1 -1
View File
@@ -90,7 +90,7 @@ their issue closes.
| `manual_transition_journal_audit.sh` | dev-tool | Journal + metrics + log audit for manual transition jobs | — |
| `manual_transition_mixed_rollout_matrix.sh` | dev-tool | Matrix generator for mixed-version rollout phases | — |
| `manual_transition_mixed_rollout_runbook.sh` | dev-tool | Reusable mixed-version rollout runbook generator (external run) | — |
| `manual_transition_mixed_version_docker_harness.sh` | dev-tool | Dedicated #1508 Docker harness for old/new manual-transition rollout evidence | `test_manual_transition_runbooks.sh` |
| `manual_transition_mixed_version_docker_harness.sh` | dev-tool | Dedicated #1508 Docker harness for old/new manual-transition rollout evidence with strict/baseline/blocked result classification | `test_manual_transition_runbooks.sh` |
| `monitor_manual_transition_ci.sh` | dev-tool | CI workflow/status watcher for manual transition follow-up monitoring | — |
| `manual_transition_soak_matrix.sh` | dev-tool | Matrix generator for nightly stress windows | — |
| `manual_transition_nightly_stress_runbook.sh` | dev-tool | Nightly stress entrypoint with failure snapshot templates | — |
@@ -25,9 +25,13 @@ COLD_IMAGE="${COLD_IMAGE:-${NEW_IMAGE}}"
BASE_PORT="${BASE_PORT:-19400}"
OUT_DIR="${OUT_DIR:-${PROJECT_ROOT}/target/manual-transition-1508-docker/$(date +%Y%m%dT%H%M%S)}"
KEEP_UP=false
ROLLBACK_NEW2_TO_OLD=true
ROLLBACK_PHASE="${ROLLBACK_PHASE:-after-terminal}"
OLD_NODE_PHASE="${OLD_NODE_PHASE:-initial}"
WAIT_TIMEOUT_SECS="${WAIT_TIMEOUT_SECS:-180}"
POLL_SECONDS="${POLL_SECONDS:-180}"
TRANSITION_WORKERS="${TRANSITION_WORKERS:-2}"
TRANSITION_QUEUE_CAPACITY="${TRANSITION_QUEUE_CAPACITY:-64}"
FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT="${FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT:-false}"
HOT_ACCESS_KEY="${HOT_ACCESS_KEY:-mvadmin}"
HOT_SECRET_KEY="${HOT_SECRET_KEY:-mvsecret}"
@@ -58,18 +62,31 @@ Options:
--object-count <n> Non-empty probe object count
--tier <name> Remote tier name
--keep-up Leave Docker services running
--no-rollback Do not replace node2 with old image after job admission
--rollback-phase <phase> Rollback node2 timing: after-terminal, in-flight, none
--old-node-phase <phase> Old node1 timing: initial, before-job
--no-rollback Alias for --rollback-phase none
-h, --help Show help
Environment:
PROJECT_NAME OLD_IMAGE NEW_IMAGE COLD_IMAGE BASE_PORT OUT_DIR KEEP_UP
HOT_ACCESS_KEY HOT_SECRET_KEY COLD_ACCESS_KEY COLD_SECRET_KEY
TIER_NAME TIER_BUCKET TIER_PREFIX JOB_BUCKET JOB_PREFIX OBJECT_COUNT
WAIT_TIMEOUT_SECS POLL_SECONDS AWS_SIGV4_SCOPE
ROLLBACK_PHASE WAIT_TIMEOUT_SECS POLL_SECONDS TRANSITION_WORKERS
TRANSITION_QUEUE_CAPACITY OLD_NODE_PHASE FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT AWS_SIGV4_SCOPE
Artifacts:
compose.yml, image inspect files, health/readiness logs, API responses,
terminal status, old-node readback, container logs, summary.env.
Result classifications:
strict_mixed_rollout_pass Real old/new images, non-empty completed transition, zero failures
baseline_tiered_storage_pass Same old/new image completed transition; useful baseline, not #1508 closure
blocked_manual_api_not_implemented Manual transition API returned 501 before job admission
blocked_manual_api_unavailable Manual transition API did not return a usable job_id
blocked_cluster_readiness_failed Docker cluster did not reach health/readiness before admission
blocked_empty_scan_or_lifecycle Job completed without lifecycle-matching transition work
blocked_manual_job_preempted_by_lifecycle_queue Lifecycle/immediate transition queued work before the job
strict_mixed_rollout_fail Mixed rollout ran but did not satisfy the strict #1508 gate
USAGE
}
@@ -111,6 +128,28 @@ parse_positive_int() {
fi
}
validate_rollback_phase() {
case "$ROLLBACK_PHASE" in
after-terminal|in-flight|none)
;;
*)
log_error "--rollback-phase must be one of: after-terminal, in-flight, none"
exit 1
;;
esac
}
validate_old_node_phase() {
case "$OLD_NODE_PHASE" in
initial|before-job)
;;
*)
log_error "--old-node-phase must be one of: initial, before-job"
exit 1
;;
esac
}
parse_args() {
while [[ $# -gt 0 ]]; do
case "$1" in
@@ -150,8 +189,16 @@ parse_args() {
KEEP_UP=true
shift
;;
--rollback-phase)
ROLLBACK_PHASE="$(arg_value "$1" "${2:-}")"
shift 2
;;
--old-node-phase)
OLD_NODE_PHASE="$(arg_value "$1" "${2:-}")"
shift 2
;;
--no-rollback)
ROLLBACK_NEW2_TO_OLD=false
ROLLBACK_PHASE=none
shift
;;
-h|--help)
@@ -198,12 +245,18 @@ cleanup() {
}
write_compose_file() {
local node1_image
COMPOSE_FILE="${OUT_DIR}/compose.yml"
node1_image="$OLD_IMAGE"
if [[ "$OLD_NODE_PHASE" == "before-job" ]]; then
node1_image="$NEW_IMAGE"
fi
cat >"$COMPOSE_FILE" <<EOF
services:
cold:
image: ${COLD_IMAGE}
hostname: cold
user: "0:0"
environment:
- RUSTFS_ADDRESS=:9000
- RUSTFS_ACCESS_KEY=${COLD_ACCESS_KEY}
@@ -222,14 +275,20 @@ services:
- rustfs-1508-net
node1:
image: ${OLD_IMAGE}
image: ${node1_image}
hostname: node1
user: "0:0"
environment: &hot-env
- RUSTFS_ADDRESS=:9000
- RUSTFS_ACCESS_KEY=${HOT_ACCESS_KEY}
- RUSTFS_SECRET_KEY=${HOT_SECRET_KEY}
- RUSTFS_VOLUMES=$(hot_volumes)
- RUSTFS_SCANNER_ENABLED=false
- RUSTFS_SCANNER_CYCLE=3600
- RUSTFS_SCANNER_START_DELAY_SECS=3600
- RUSTFS_MAX_TRANSITION_WORKERS=${TRANSITION_WORKERS}
- RUSTFS_TRANSITION_QUEUE_CAPACITY=${TRANSITION_QUEUE_CAPACITY}
- RUSTFS_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT=${FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT}
- RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true
- RUSTFS_OBS_LOGGER_LEVEL=warn
volumes:
@@ -245,6 +304,7 @@ services:
node2:
image: ${NEW_IMAGE}
hostname: node2
user: "0:0"
environment: *hot-env
volumes:
- node2_data_0:/data/rustfs0
@@ -259,6 +319,7 @@ services:
node3:
image: ${NEW_IMAGE}
hostname: node3
user: "0:0"
environment: *hot-env
volumes:
- node3_data_0:/data/rustfs0
@@ -273,6 +334,7 @@ services:
node4:
image: ${NEW_IMAGE}
hostname: node4
user: "0:0"
environment: *hot-env
volumes:
- node4_data_0:/data/rustfs0
@@ -433,16 +495,19 @@ seed_objects() {
}
start_transition_job() {
local query response job_id
local query response job_id run_http_code
query="bucket=$(url_encode "$JOB_BUCKET")"
query="${query}&prefix=$(url_encode "${JOB_PREFIX}/")"
query="${query}&tier=$(url_encode "$TIER_NAME")"
query="${query}&dryRun=false&maxObjects=${OBJECT_COUNT}&mode=async"
response="$(curl_hot POST "$(hot_endpoint 2)/rustfs/admin/v3/ilm/transition/run?${query}")"
printf '%s\n' "$response" >"${OUT_DIR}/run-response.json"
curl_hot POST "$(hot_endpoint 2)/rustfs/admin/v3/ilm/transition/run?${query}" \
-o "${OUT_DIR}/run-response.json" \
-w "%{http_code}\n" >"${OUT_DIR}/run-response.http_code" || true
response="$(cat "${OUT_DIR}/run-response.json" 2>/dev/null || true)"
run_http_code="$(cat "${OUT_DIR}/run-response.http_code" 2>/dev/null || true)"
job_id="$(printf '%s' "$response" | jq -r '.job_id // empty')"
if [[ -z "$job_id" ]]; then
log_error "manual transition response omitted job_id"
log_warn "manual transition response omitted job_id, http_code=${run_http_code:-unknown}"
return 1
fi
printf '%s\n' "$job_id" >"${OUT_DIR}/job-id.txt"
@@ -451,19 +516,25 @@ start_transition_job() {
replace_node2_with_old_image() {
local network
network="$(network_name)"
log_info "Replacing node2 with old image ${OLD_IMAGE} for in-flight rollback readback"
log_info "Replacing node2 with old image ${OLD_IMAGE} for ${ROLLBACK_PHASE} rollback readback"
compose stop node2 >/dev/null
compose rm -f node2 >/dev/null
docker run -d \
--name "${PROJECT_NAME}-node2-rollback-old" \
--network "$network" \
--network-alias node2 \
--user 0:0 \
-p "$((BASE_PORT + 2)):9000" \
-e RUSTFS_ADDRESS=:9000 \
-e RUSTFS_ACCESS_KEY="$HOT_ACCESS_KEY" \
-e RUSTFS_SECRET_KEY="$HOT_SECRET_KEY" \
-e RUSTFS_VOLUMES="$(hot_volumes)" \
-e RUSTFS_SCANNER_ENABLED=false \
-e RUSTFS_SCANNER_CYCLE=3600 \
-e RUSTFS_SCANNER_START_DELAY_SECS=3600 \
-e RUSTFS_MAX_TRANSITION_WORKERS="$TRANSITION_WORKERS" \
-e RUSTFS_TRANSITION_QUEUE_CAPACITY="$TRANSITION_QUEUE_CAPACITY" \
-e RUSTFS_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT="$FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT" \
-e RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true \
-e RUSTFS_OBS_LOGGER_LEVEL=warn \
-v "${PROJECT_NAME}_node2_data_0:/data/rustfs0" \
@@ -473,6 +544,37 @@ replace_node2_with_old_image() {
"$OLD_IMAGE" >/dev/null
}
replace_node1_with_old_image() {
local network
network="$(network_name)"
log_info "Replacing node1 with old image ${OLD_IMAGE} before manual transition job"
compose stop node1 >/dev/null
compose rm -f node1 >/dev/null
docker run -d \
--name "${PROJECT_NAME}-node1-before-job-old" \
--network "$network" \
--network-alias node1 \
--user 0:0 \
-p "$((BASE_PORT + 1)):9000" \
-e RUSTFS_ADDRESS=:9000 \
-e RUSTFS_ACCESS_KEY="$HOT_ACCESS_KEY" \
-e RUSTFS_SECRET_KEY="$HOT_SECRET_KEY" \
-e RUSTFS_VOLUMES="$(hot_volumes)" \
-e RUSTFS_SCANNER_ENABLED=false \
-e RUSTFS_SCANNER_CYCLE=3600 \
-e RUSTFS_SCANNER_START_DELAY_SECS=3600 \
-e RUSTFS_MAX_TRANSITION_WORKERS="$TRANSITION_WORKERS" \
-e RUSTFS_TRANSITION_QUEUE_CAPACITY="$TRANSITION_QUEUE_CAPACITY" \
-e RUSTFS_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT="$FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT" \
-e RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true \
-e RUSTFS_OBS_LOGGER_LEVEL=warn \
-v "${PROJECT_NAME}_node1_data_0:/data/rustfs0" \
-v "${PROJECT_NAME}_node1_data_1:/data/rustfs1" \
-v "${PROJECT_NAME}_node1_data_2:/data/rustfs2" \
-v "${PROJECT_NAME}_node1_data_3:/data/rustfs3" \
"$OLD_IMAGE" >/dev/null
}
poll_terminal_status() {
local job_id="$1"
local status_url status_json terminal_state
@@ -512,6 +614,85 @@ head_probe() {
-w "%{http_code}\n" >"${OUT_DIR}/head-object.http_code" || true
}
image_id() {
local file="$1"
jq -r '.[0].Id // ""' "$file" 2>/dev/null || true
}
is_compat_readback_code() {
case "$1" in
200|501)
return 0
;;
*)
return 1
;;
esac
}
classify_result() {
local terminal_state="$1"
local transition_completed="$2"
local transition_failed="$3"
local tier_failure="$4"
local old_code="$5"
local rollback_code="$6"
local run_http_code="$7"
local lifecycle_config_found="$8"
local scanned="$9"
local eligible="${10}"
local skipped_already_transitioned="${11}"
local skipped_already_in_flight="${12}"
local old_image_id new_image_id images_are_mixed readback_ok
if [[ -f "${OUT_DIR}/readiness-failed" ]]; then
printf 'blocked_cluster_readiness_failed\n'
return
fi
old_image_id="$(image_id "${OUT_DIR}/old-image.inspect.json")"
new_image_id="$(image_id "${OUT_DIR}/new-image.inspect.json")"
images_are_mixed=false
if [[ "$OLD_IMAGE" != "$NEW_IMAGE" && -n "$old_image_id" && -n "$new_image_id" && "$old_image_id" != "$new_image_id" ]]; then
images_are_mixed=true
fi
readback_ok=false
if is_compat_readback_code "$old_code"; then
if [[ "$ROLLBACK_PHASE" == "none" ]] || is_compat_readback_code "$rollback_code"; then
readback_ok=true
fi
fi
if [[ "$run_http_code" == "501" ]]; then
printf 'blocked_manual_api_not_implemented\n'
return
fi
if [[ -z "$(cat "${OUT_DIR}/job-id.txt" 2>/dev/null || true)" ]]; then
printf 'blocked_manual_api_unavailable\n'
return
fi
if [[ "$terminal_state" == "completed" && ( "$lifecycle_config_found" != "true" || "$scanned" == "0" || "$eligible" == "0" || "$transition_completed" == "0" ) ]]; then
printf 'blocked_empty_scan_or_lifecycle\n'
return
fi
if [[ "$transition_completed" == "0" && "$eligible" != "0" && "$transition_failed" == "0" && "$tier_failure" == "0" ]]; then
if [[ "$skipped_already_transitioned" != "0" || "$skipped_already_in_flight" != "0" ]]; then
printf 'blocked_manual_job_preempted_by_lifecycle_queue\n'
return
fi
fi
if [[ "$terminal_state" == "completed" && "$transition_completed" != "0" && "$transition_failed" == "0" && "$tier_failure" == "0" ]]; then
if [[ "$images_are_mixed" == "true" && "$readback_ok" == "true" ]]; then
printf 'strict_mixed_rollout_pass\n'
return
fi
printf 'baseline_tiered_storage_pass\n'
return
fi
printf 'strict_mixed_rollout_fail\n'
}
collect_logs() {
local service
if [[ -z "$COMPOSE_FILE" || ! -f "$COMPOSE_FILE" ]]; then
@@ -520,37 +701,63 @@ collect_logs() {
for service in cold node1 node2 node3 node4; do
compose logs --no-color "$service" >"${OUT_DIR}/${service}.log" 2>/dev/null || true
done
docker logs "${PROJECT_NAME}-node1-before-job-old" >"${OUT_DIR}/node1-before-job-old.log" 2>/dev/null || true
docker logs "${PROJECT_NAME}-node2-rollback-old" >"${OUT_DIR}/node2-rollback-old.log" 2>/dev/null || true
}
summarize() {
local terminal_state transition_completed tier_failure transition_failed old_code rollback_code
local terminal_state transition_completed tier_failure transition_failed old_code rollback_code run_http_code lifecycle_config_found scanned eligible skipped_already_transitioned skipped_already_in_flight queue_queued queue_active result_classification
terminal_state="$(cat "${OUT_DIR}/terminal-state.txt" 2>/dev/null || true)"
transition_completed="$(jq -r '.report.transition_completed // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
tier_failure="$(jq -r '.report.tier_failure // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
transition_failed="$(jq -r '.report.transition_failed // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
old_code="$(cat "${OUT_DIR}/old-node-status.http_code" 2>/dev/null || true)"
rollback_code="$(cat "${OUT_DIR}/rollback-node2-status.http_code" 2>/dev/null || true)"
run_http_code="$(cat "${OUT_DIR}/run-response.http_code" 2>/dev/null || true)"
lifecycle_config_found="$(jq -r '.report.lifecycle_config_found // false' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf 'false')"
scanned="$(jq -r '.report.scanned // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
eligible="$(jq -r '.report.eligible // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
skipped_already_transitioned="$(jq -r '.report.skipped_already_transitioned // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
skipped_already_in_flight="$(jq -r '.report.skipped_already_in_flight // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
queue_queued="$(jq -r '.queue_snapshot.queued // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
queue_active="$(jq -r '.queue_snapshot.active // 0' "${OUT_DIR}/status-terminal.json" 2>/dev/null || printf '0')"
result_classification="$(classify_result "$terminal_state" "$transition_completed" "$transition_failed" "$tier_failure" "$old_code" "$rollback_code" "$run_http_code" "$lifecycle_config_found" "$scanned" "$eligible" "$skipped_already_transitioned" "$skipped_already_in_flight")"
cat >"${OUT_DIR}/summary.env" <<EOF
project_name=${PROJECT_NAME}
old_image=${OLD_IMAGE}
new_image=${NEW_IMAGE}
cold_image=${COLD_IMAGE}
old_image_id=$(image_id "${OUT_DIR}/old-image.inspect.json")
new_image_id=$(image_id "${OUT_DIR}/new-image.inspect.json")
base_port=${BASE_PORT}
tier=${TIER_NAME}
job_bucket=${JOB_BUCKET}
job_prefix=${JOB_PREFIX}
object_count=${OBJECT_COUNT}
transition_workers=${TRANSITION_WORKERS}
transition_queue_capacity=${TRANSITION_QUEUE_CAPACITY}
force_immediate_transition_enqueue_timeout=${FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT}
run_http_code=${run_http_code}
terminal_state=${terminal_state}
lifecycle_config_found=${lifecycle_config_found}
scanned=${scanned}
eligible=${eligible}
skipped_already_transitioned=${skipped_already_transitioned}
skipped_already_in_flight=${skipped_already_in_flight}
transition_completed=${transition_completed}
transition_failed=${transition_failed}
tier_failure=${tier_failure}
queue_queued=${queue_queued}
queue_active=${queue_active}
old_node_status_http_code=${old_code}
rollback_node2_status_http_code=${rollback_code}
rollback_new2_to_old=${ROLLBACK_NEW2_TO_OLD}
rollback_phase=${ROLLBACK_PHASE}
old_node_phase=${OLD_NODE_PHASE}
rollback_new2_to_old=$([[ "$ROLLBACK_PHASE" == "none" ]] && printf 'false' || printf 'true')
result_classification=${result_classification}
EOF
cat "${OUT_DIR}/summary.env"
if [[ "$terminal_state" != "completed" || "$transition_completed" == "0" || "$transition_failed" != "0" || "$tier_failure" != "0" ]]; then
if [[ "$result_classification" != "strict_mixed_rollout_pass" ]]; then
log_error "strict #1508 mixed-version transition gate failed; see ${OUT_DIR}"
return 1
fi
@@ -560,6 +767,8 @@ main() {
parse_args "$@"
parse_positive_int "--base-port" "$BASE_PORT"
parse_positive_int "--object-count" "$OBJECT_COUNT"
validate_rollback_phase
validate_old_node_phase
require_cmd docker
require_cmd curl
require_cmd jq
@@ -574,20 +783,37 @@ main() {
docker image inspect "$COLD_IMAGE" >"${OUT_DIR}/cold-image.inspect.json"
compose up -d
wait_cluster_ready
if ! wait_cluster_ready; then
printf 'cluster readiness failed before manual transition admission\n' >"${OUT_DIR}/readiness-failed"
collect_logs
summarize
return 1
fi
curl_cold PUT "$(cold_endpoint)/${TIER_BUCKET}" -o "${OUT_DIR}/create-cold-bucket.response" -w "%{http_code}\n" >"${OUT_DIR}/create-cold-bucket.http_code"
create_bucket "$(hot_endpoint 2)" "$JOB_BUCKET" "${OUT_DIR}/create-hot-bucket"
add_tier
put_lifecycle
seed_objects
start_transition_job
if [[ "$ROLLBACK_NEW2_TO_OLD" == "true" ]]; then
replace_node2_with_old_image
put_lifecycle
if [[ "$OLD_NODE_PHASE" == "before-job" ]]; then
replace_node1_with_old_image
wait_http_ok "$(hot_endpoint 1)/health" "node1-old-live" || true
wait_http_ok "$(hot_endpoint 1)/health/ready" "node1-old-ready" || true
fi
if start_transition_job; then
if [[ "$ROLLBACK_PHASE" == "in-flight" ]]; then
replace_node2_with_old_image
fi
poll_terminal_status "$(cat "${OUT_DIR}/job-id.txt")" || true
if [[ "$ROLLBACK_PHASE" == "after-terminal" ]]; then
replace_node2_with_old_image
wait_http_ok "$(hot_endpoint 2)/health" "rollback-node2-live" || true
fi
capture_old_node_readback "$(cat "${OUT_DIR}/job-id.txt")"
head_probe
else
log_warn "Skipping terminal polling because no manual transition job was admitted"
fi
poll_terminal_status "$(cat "${OUT_DIR}/job-id.txt")"
capture_old_node_readback "$(cat "${OUT_DIR}/job-id.txt")"
head_probe
collect_logs
summarize
}
+2 -2
View File
@@ -317,7 +317,7 @@ write_blackbox_matrix() {
printf 'quick\theal degraded erasure disk rebuild\tblack-box\tcargo test --package e2e_test heal_erasure_disk_rebuild_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'quick\tnamespace lock quorum under EC ops\tblack-box\tcargo test --package e2e_test namespace_lock_quorum_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'full\tlegacy bitrot read fixture restore\tfixture\tcargo test -p rustfs-ecstore --test legacy_bitrot_read_test -- --nocapture\tRUSTFS_LEGACY_TEST_ROOT,RUSTFS_LEGACY_TEST_DISK\t%s\n' "$legacy_status"
printf 'full\tMinIO generated encrypted read and negative restore fixture\tfixture\tcargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored --nocapture\tRUSTFS_MINIO_FIXTURE_ROOT,RUSTFS_MINIO_STATIC_KMS_KEY_B64\t%s\n' "$minio_status"
printf 'full\tMinIO generated encrypted read and negative restore fixture\tfixture\tcargo test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored --nocapture\tRUSTFS_MINIO_FIXTURE_ROOT,RUSTFS_MINIO_STATIC_KMS_KEY_B64\t%s\n' "$minio_status"
printf 'full\tS3 multipart range versioning delete subset\tblack-box\tenv TESTEXPR=\"multipart or range or versioning or delete\" DEPLOY_MODE=build MAXFAIL=0 ./scripts/s3-tests/run.sh\tnone\t%s\n' "$s3_status"
printf 'destructive\tdistributed cluster concurrency\tblack-box\tcargo test --package e2e_test cluster_concurrency_test -- --nocapture\tnone\t%s\n' "$destructive_status"
printf 'destructive\tstale multipart cleanup cluster\tblack-box\tcargo test --package e2e_test stale_multipart_cleanup_cluster_test -- --nocapture\tnone\t%s\n' "$destructive_status"
@@ -377,7 +377,7 @@ run_fixture_steps() {
if fixture_available; then
run_step "ecstore-minio-generated-read-fixture" \
cargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored --nocapture
cargo test -p rustfs-ecstore --features rio-v2 --test minio_generated_read_test -- --ignored --nocapture
elif [[ "$REQUIRE_FIXTURES" == "true" ]]; then
echo "ERROR: $(minio_fixture_missing_reason)" >&2
exit 1
@@ -208,7 +208,16 @@ bash "$MIXED_DOCKER_HARNESS" --help >/tmp/manual_transition_mixed_version_docker
rg -q "mixed_version_docker_harness" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q -- "--old-image" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q -- "--new-image" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q -- "--rollback-phase" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q -- "--old-node-phase" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q -- "--no-rollback" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "strict_mixed_rollout_pass" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "baseline_tiered_storage_pass" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "blocked_manual_api_not_implemented" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "blocked_cluster_readiness_failed" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "blocked_manual_job_preempted_by_lifecycle_queue" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "OLD_NODE_PHASE" /tmp/manual_transition_mixed_version_docker_harness.help
rg -q "FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT" /tmp/manual_transition_mixed_version_docker_harness.help
if bash "$FAILURE_SAMPLES" --endpoint http://127.0.0.1:9000 --sample >/tmp/manual_transition_failure_samples.err 2>&1; then
echo "failure samples script should fail when --sample has no value" >&2
exit 1