Compare commits

..

5 Commits

Author SHA1 Message Date
Zhengchao An 23a2c7d776 test(kms): stabilize Vault failover validation (#6385)
* test(kms): bound Vault failover progress wait

* ci(nightly): honor manual dispatch ref

* test(kms): preserve Vault worker failures

* test(kms): validate Vault circuit recovery
2026-08-23 12:32:11 +08:00
唐小鸭 5f72209446 fix(ecstore): keep unknown-size sentinel in create_bitrot_writer (#6380)
SSE and compression wrap the payload so its length is unknown and
advertise HashReader::SIZE_PRESERVE_LAYER (-1). Every layer preserved
that sentinel except create_bitrot_writer, which clamped it to 0 before
calling DiskAPI::create_file. RemoteDisk forwards that size verbatim in
the put_file_stream query, so remote peers were told the body was empty.

Since the authenticated put-file trailer (#5868) the receiver used the
declared size to split body from trailer, turning the clamp into a fatal
"auth trailer has trailing data" failure for every SSE PUT on multi-node
deployments (rc.2). #6320 relaxed the receiver to only trust size > 0;
this change fixes the sender so the sentinel survives end to end and the
wire no longer conflates empty objects with unknown-length streams.

Refs #6331
2026-08-23 12:29:52 +08:00
Zhengchao An b6ba89d9e4 docs(testing): document CI gate matrix (#6412) 2026-08-23 12:09:06 +08:00
houseme 648d5166e2 feat(allocator): replace mimalloc/libmimalloc-sys with rustfs-mimalloc/rustfs-mimalloc-sys (#6404)
Replace the upstream xonatius/mimalloc_rust.git fork (mimalloc + libmimalloc-sys)
with the published rustfs-mimalloc (v0.5.0) and rustfs-mimalloc-sys (v0.5.0) crates
from crates.io.

The new crates are based on mimalloc V3 (v3.5.0) and provide:
- MiMalloc global allocator with safe API (collect, stats_json, process_info)
- Heap management and arena operations (heap module)
- Full FFI bindings to mimalloc V3

Changes:
- Workspace deps: mimalloc + libmimalloc-sys (git) → rustfs-mimalloc + rustfs-mimalloc-sys (crates.io)
- allocator_reclaim.rs: libmimalloc_sys::mi_collect → rustfs_mimalloc::MiMalloc::collect
- memory_observability.rs: raw FFI mi_stats_get_json → MiMalloc::stats_json()
- main.rs: heap ownership tests use Heap::contains() (V3 API)
- deny.toml: remove xonatius/mimalloc_rust.git from allow-git

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-23 12:07:25 +08:00
houseme 84eb5aebef fix(ecstore): remove inline write debug noise (#6408)
* fix(ecstore): remove inline write debug noise

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

* fix(ecstore): satisfy warning-as-error lints

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

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-23 12:07:20 +08:00
23 changed files with 529 additions and 1959 deletions
+3 -6
View File
@@ -39,11 +39,10 @@ jobs:
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -89,11 +88,10 @@ jobs:
# either casing.
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -178,11 +176,10 @@ jobs:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
+2 -1
View File
@@ -30,7 +30,8 @@ make build-docker BUILD_OS=ubuntu22.04
- Crate membership: `Cargo.toml` `[workspace].members`
- Architecture, layering, crate map: [ARCHITECTURE.md](ARCHITECTURE.md)
- Migration guardrails & readiness contracts: [docs/architecture/](docs/architecture/README.md)
- CI gates: `.github/workflows/ci.yml` (source of truth; never copy its steps into docs)
- CI workflow steps: `.github/workflows/`; event, timeout, and required-status
matrix: [docs/testing/ci-gates.md](docs/testing/ci-gates.md)
- Test-layer taxonomy, per-layer entry commands, serial/nextest rules, flake
policy: [docs/testing/README.md](docs/testing/README.md)
- Tier/ILM transition debugging (xl.meta inspection, versionId tracing):
+2
View File
@@ -70,6 +70,8 @@ make pre-pr
> For the full test-layer taxonomy (unit / ecstore black-box / e2e / s3s-e2e / S3 compatibility / chaos / fuzz / bench), each layer's entry command, the naming conventions the migration gate depends on, and the serial/nextest rules, see [docs/testing/README.md](docs/testing/README.md).
> For the event, timeout, required-status, and local reproduction matrix, see [docs/testing/ci-gates.md](docs/testing/ci-gates.md).
### 🔒 Automated Pre-commit Hooks
#### What `make pre-commit` and `make pre-pr` actually run
Generated
+22 -27
View File
@@ -1858,9 +1858,9 @@ dependencies = [
[[package]]
name = "cc"
version = "1.4.3"
version = "1.4.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d"
checksum = "0ad534f4357a5264cce5019c989cf66a4f0dc4e0d1b1d15f8aacec0ff7360273"
dependencies = [
"find-msvc-tools",
"jobserver",
@@ -2522,12 +2522,6 @@ dependencies = [
"subtle",
]
[[package]]
name = "cty"
version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b365fabc795046672053e29c954733ec3b05e4be654ab130fe8f1f94d7051f35"
[[package]]
name = "curve25519-dalek"
version = "4.1.3"
@@ -5988,15 +5982,6 @@ version = "0.2.16"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981"
[[package]]
name = "libmimalloc-sys"
version = "0.1.49"
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
dependencies = [
"cc",
"cty",
]
[[package]]
name = "libredox"
version = "0.1.20"
@@ -6397,14 +6382,6 @@ dependencies = [
"synstructure 0.13.2",
]
[[package]]
name = "mimalloc"
version = "0.1.52"
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
dependencies = [
"libmimalloc-sys",
]
[[package]]
name = "mime"
version = "0.3.17"
@@ -9162,13 +9139,11 @@ dependencies = [
"insta",
"jiff",
"libc",
"libmimalloc-sys",
"libsystemd",
"matchit 0.9.2",
"md-5 0.11.0",
"metrics",
"metrics-util",
"mimalloc",
"mime_guess",
"opentelemetry",
"opentelemetry_sdk",
@@ -9204,6 +9179,8 @@ dependencies = [
"rustfs-lock",
"rustfs-log-analyzer",
"rustfs-madmin",
"rustfs-mimalloc",
"rustfs-mimalloc-sys",
"rustfs-notify",
"rustfs-object-capacity",
"rustfs-object-data-cache",
@@ -9875,6 +9852,24 @@ dependencies = [
"tokio",
]
[[package]]
name = "rustfs-mimalloc"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a406f4aa07084301d485beec873af6dccc8e3f8762da244743df92038b1db1a6"
dependencies = [
"rustfs-mimalloc-sys",
]
[[package]]
name = "rustfs-mimalloc-sys"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c3051b819175f58445d4c369a72f0ab88149f3885ba8bea2aff3be01f53fe7cd"
dependencies = [
"cc",
]
[[package]]
name = "rustfs-notify"
version = "1.0.0-rc.3"
+2 -2
View File
@@ -350,8 +350,8 @@ russh-sftp = "2.4.0"
dav-server = "0.11.0"
# Performance Analysis and Memory Profiling
mimalloc = { version = "0.1.52", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11" }
libmimalloc-sys = { version = "0.1.49", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11", features = ["extended"] }
rustfs-mimalloc = { version = "0.5.0" }
rustfs-mimalloc-sys = { version = "0.5.0" }
hotpath = { version = "0.23.3", default-features = false }
# Snapshot testing for output format regression detection
insta = { version = "1.48" }
+38 -6
View File
@@ -784,6 +784,24 @@ pub(crate) fn create_deferred_bitrot_reader_with_stripe_handle(
///
/// # Returns
/// A Result containing the BitrotWriterWrapper or an error
/// Size hint handed to `DiskAPI::create_file` for a bitrot-wrapped shard.
///
/// A known length is grown by one checksum per shard so the on-disk file size
/// matches what the bitrot writer emits. A negative length is the
/// unknown-size sentinel (`HashReader::SIZE_PRESERVE_LAYER`, used by SSE and
/// compression) and must be preserved: `RemoteDisk::create_file` forwards it
/// in the `put_file_stream` query, and the receiver only treats `size > 0` as
/// a fixed body length when locating the authenticated trailer. Clamping it
/// to `0` would claim an empty body and misframe the stream. `0` stays `0`
/// because a genuinely empty object still means an empty body.
fn bitrot_create_file_size(length: i64, shard_size: usize, checksum_algo: &HashAlgorithm) -> i64 {
if length <= 0 {
return length;
}
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
}
pub async fn create_bitrot_writer(
is_inline_buffer: bool,
disk: Option<&DiskStore>,
@@ -796,12 +814,7 @@ pub async fn create_bitrot_writer(
let writer = if is_inline_buffer {
CustomWriter::new_inline_buffer()
} else if let Some(disk) = disk {
let length = if length > 0 {
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
} else {
0
};
let length = bitrot_create_file_size(length, shard_size, &checksum_algo);
let file = disk.create_file("", volume, path, length).await?;
#[cfg(feature = "hotpath")]
@@ -820,6 +833,25 @@ mod tests {
use rustfs_rio::ChunkReader;
use std::collections::VecDeque;
#[test]
fn bitrot_create_file_size_grows_known_length_by_checksums() {
// 10 bytes over 4-byte shards = 3 shards, each followed by a 32-byte hash.
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::HighwayHash256), 10 + 3 * 32);
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::None), 10);
}
#[test]
fn bitrot_create_file_size_keeps_empty_and_unknown_distinct() {
assert_eq!(bitrot_create_file_size(0, 4, &HashAlgorithm::HighwayHash256), 0);
// SSE/compression streams advertise SIZE_PRESERVE_LAYER (-1); the remote
// put_file_stream receiver relies on a non-positive size to parse the auth
// trailer from the stream tail, so the sentinel must survive untouched.
assert_eq!(
bitrot_create_file_size(rustfs_rio::HashReader::SIZE_PRESERVE_LAYER, 4, &HashAlgorithm::HighwayHash256),
rustfs_rio::HashReader::SIZE_PRESERVE_LAYER
);
}
struct TestChunkReader {
chunks: VecDeque<Bytes>,
}
+1 -14
View File
@@ -2124,26 +2124,13 @@ impl SetDisks {
let put_object_size = known_put_object_storage_size(data.size());
let shard_file_size_raw = erasure.shard_file_size(put_object_size);
let is_inline_buffer =
storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
let is_inline_buffer = storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled();
let shard_file_size = shard_file_size_raw;
let shard_size = erasure.shard_size();
let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size);
let direct_inline_commit = matches!(write_path, SmallWritePath::Inline);
{
use std::io::Write;
let msg = format!(
"INLINE_DEBUG: bucket={} obj={} size={} shard_fs={} ds={} bs={} inline={} direct={} path={} iblock={} ver={}\n",
bucket, object, put_object_size, shard_file_size_raw, erasure.data_shards, fi.erasure.block_size,
is_inline_buffer, direct_inline_commit, write_path.metric_label(), storage_class_config.inline_block(), opts.versioned
);
if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open("/tmp/rustfs_inline_debug.log") {
let _ = f.write_all(msg.as_bytes());
}
let _ = std::io::stderr().write_all(msg.as_bytes());
}
rustfs_io_metrics::record_put_object_path(write_path.metric_label());
let writer_setup_stage_start = collect_stage_timing.then(Instant::now);
let (mut writers, errors) = if direct_inline_commit {
+86 -25
View File
@@ -16,14 +16,14 @@
//!
//! `scripts/test/vault_ha_kms_live.sh` owns the official Vault containers and
//! kills the active node while this test continuously decrypts through a
//! surviving standby. KV2 and Transit requests must remain successful, use a
//! bounded number of attempts, and leave the circuit and in-flight gauges at
//! zero after a new leader is elected.
//! surviving standby. KV2 and Transit must recover after the bounded circuit
//! interval, use a bounded number of attempts, and leave the circuit and
//! in-flight gauges at zero after a new leader is elected.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use metrics_util::MetricKind;
@@ -43,6 +43,11 @@ const OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
const IN_FLIGHT: &str = "rustfs_kms_backend_in_flight";
const CIRCUIT_OPEN: &str = "rustfs_kms_backend_circuit_open";
const MAX_ATTEMPTS: u32 = 10;
const ATTEMPT_TIMEOUT: Duration = Duration::from_secs(2);
const HEALTHY_PROGRESS_TIMEOUT: Duration = Duration::from_secs(20);
// The circuit remains open for 30s after five failed attempts.
const POST_FAILOVER_PROGRESS_TIMEOUT: Duration = Duration::from_secs(35);
const FAILOVER_ERROR_POLL_INTERVAL: Duration = Duration::from_millis(100);
type MetricEntry = (
metrics_util::CompositeKey,
@@ -64,7 +69,7 @@ fn config(backend: KmsBackend, backend_config: BackendConfig) -> KmsConfig {
backend,
backend_config,
allow_insecure_dev_defaults: true,
timeout: Duration::from_secs(2),
timeout: ATTEMPT_TIMEOUT,
retry_attempts: MAX_ATTEMPTS,
enable_cache: false,
..KmsConfig::default()
@@ -164,14 +169,31 @@ fn retryable_failures(snapshot: &[MetricEntry], operation: &str) -> u64 {
.sum()
}
async fn wait_for_count(counter: &AtomicU64, minimum: u64, description: &str) {
tokio::time::timeout(Duration::from_secs(20), async {
async fn wait_for_count(
counter: &AtomicU64,
failure: &Mutex<Option<String>>,
minimum: u64,
description: &str,
timeout: Duration,
) {
tokio::time::timeout(timeout, async {
while counter.load(Ordering::SeqCst) < minimum {
if let Some(error) = failure.lock().expect("decrypt failure lock poisoned").as_ref() {
panic!(
"{description} worker failed after {} successful decrypts: {error}",
counter.load(Ordering::SeqCst)
);
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
.unwrap_or_else(|_| {
panic!(
"timed out after {timeout:?} waiting for {description}: completed {}, expected {minimum}",
counter.load(Ordering::SeqCst)
)
});
}
async fn wait_for_file(path: &Path, description: &str) {
@@ -189,7 +211,8 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
request: DecryptRequest,
expected: Vec<u8>,
completed: Arc<AtomicU64>,
failed: Arc<AtomicBool>,
allow_failover_errors: Arc<AtomicBool>,
failure: Arc<Mutex<Option<String>>>,
stop: CancellationToken,
) {
while !stop.is_cancelled() {
@@ -197,8 +220,18 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
Ok(response) if response.plaintext == expected => {
completed.fetch_add(1, Ordering::SeqCst);
}
Ok(_) | Err(_) => {
failed.store(true, Ordering::SeqCst);
Ok(_) => {
*failure.lock().expect("decrypt failure lock poisoned") =
Some("decrypt returned unexpected plaintext".to_string());
return;
}
Err(rustfs_kms::KmsError::BackendError { .. } | rustfs_kms::KmsError::OperationTimedOut { .. })
if allow_failover_errors.load(Ordering::SeqCst) =>
{
tokio::time::sleep(FAILOVER_ERROR_POLL_INTERVAL).await;
}
Err(error) => {
*failure.lock().expect("decrypt failure lock poisoned") = Some(error.to_string());
return;
}
}
@@ -296,7 +329,9 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
);
let stop = CancellationToken::new();
let failed = Arc::new(AtomicBool::new(false));
let allow_failover_errors = Arc::new(AtomicBool::new(false));
let kv2_failure = Arc::new(Mutex::new(None));
let transit_failure = Arc::new(Mutex::new(None));
let kv2_completed = Arc::new(AtomicU64::new(0));
let transit_completed = Arc::new(AtomicU64::new(0));
let kv2_worker = tokio::spawn(decrypt_loop(
@@ -304,7 +339,8 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
kv2_request,
kv2_data_key.plaintext_key,
Arc::clone(&kv2_completed),
Arc::clone(&failed),
Arc::clone(&allow_failover_errors),
Arc::clone(&kv2_failure),
stop.clone(),
));
let transit_worker = tokio::spawn(decrypt_loop(
@@ -312,12 +348,21 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
transit_request,
transit_data_key.plaintext_key,
Arc::clone(&transit_completed),
Arc::clone(&failed),
Arc::clone(&allow_failover_errors),
Arc::clone(&transit_failure),
stop.clone(),
));
wait_for_count(&kv2_completed, 2, "two healthy KV2 decrypts").await;
wait_for_count(&transit_completed, 2, "two healthy Transit decrypts").await;
wait_for_count(&kv2_completed, &kv2_failure, 2, "two healthy KV2 decrypts", HEALTHY_PROGRESS_TIMEOUT).await;
wait_for_count(
&transit_completed,
&transit_failure,
2,
"two healthy Transit decrypts",
HEALTHY_PROGRESS_TIMEOUT,
)
.await;
allow_failover_errors.store(true, Ordering::SeqCst);
std::fs::write(&marker, b"ready").expect("publish failover readiness marker");
wait_for_file(&elected, "the replacement Vault leader").await;
@@ -326,18 +371,39 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
let kv2_after_election = kv2_completed.load(Ordering::SeqCst) + 2;
let transit_after_election = transit_completed.load(Ordering::SeqCst) + 2;
wait_for_count(&kv2_completed, kv2_after_election, "post-failover KV2 decrypts").await;
wait_for_count(&transit_completed, transit_after_election, "post-failover Transit decrypts").await;
wait_for_count(
&kv2_completed,
&kv2_failure,
kv2_after_election,
"post-failover KV2 decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(
&transit_completed,
&transit_failure,
transit_after_election,
"post-failover Transit decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
stop.cancel();
kv2_worker.await.expect("KV2 decrypt worker must join");
transit_worker.await.expect("Transit decrypt worker must join");
assert!(!failed.load(Ordering::SeqCst), "no decrypt may fail or return different plaintext");
assert!(
kv2_failure.lock().expect("KV2 failure lock poisoned").is_none(),
"no KV2 decrypt may fail or return different plaintext"
);
assert!(
transit_failure.lock().expect("Transit failure lock poisoned").is_none(),
"no Transit decrypt may fail or return different plaintext"
);
}
#[test]
#[ignore = "requires a real three-node Vault Raft cluster; run scripts/test/vault_ha_kms_live.sh"]
fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
@@ -349,11 +415,6 @@ fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
});
let snapshot = snapshotter.snapshot().into_vec();
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "circuit_open")]),
0,
"a bounded leader election must not open the circuit"
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]),
0,
-3
View File
@@ -43,9 +43,6 @@ allow-git = [
# RustFS fork carrying presigned expiry and constant-time authentication fixes.
# owner: rustfs-maintainers review: 2026-10
"https://github.com/rustfs/s3s.git",
# MiMalloc fork pinned for hotpath allocation counting support.
# owner: houseme review: 2026-10
"https://github.com/xonatius/mimalloc_rust.git",
]
[bans]
+149
View File
@@ -0,0 +1,149 @@
# CI gate matrix
This file is the source of truth for which validation runs on each event, its
configured wall-clock budget, and whether it can block a merge. Test taxonomy,
naming, and nextest serialization rules remain in [README.md](README.md); e2e
membership and counts remain in
[e2e-suite-inventory.md](e2e-suite-inventory.md).
The distinction between **required** and **report-only** is load-bearing:
a failing job blocks a merge only when its exact check name is present in the
live `main` ruleset. A workflow name, a `merge_group` trigger, or a red PR check
does not make a job required by itself.
## Required merge checks
The live `main` ruleset (`6436880`) currently requires exactly these contexts:
| Required context | Producer | Validation |
|---|---|---|
| `CLA Check` | `.github/workflows/cla.yml` | Contributor agreement |
| `Quick Checks` | `.github/workflows/ci.yml` | Formatting and repository guard scripts |
| `Test and Lint` | `.github/workflows/ci.yml` | Clippy, workspace nextest excluding `e2e_test`, doctests, and migration proofs |
For pull requests limited to the paths excluded by the main CI workflow,
`.github/workflows/ci-docs-only.yml` reports `Quick Checks` and
`Test and Lint` under the same names. It runs the real quick checks and the
planning-document guard; it does not claim that Rust compilation or runtime
tests ran. Despite the workflow name, these paths also include selected deploy,
workflow, and lock files.
Verify the live rule rather than trusting this snapshot before changing merge
policy:
```bash
gh api repos/rustfs/rustfs/rulesets/6436880 \
--jq '.rules[] | select(.type == "required_status_checks") | .parameters'
```
The ruleset currently has `strict_required_status_checks_policy=false`.
`Continuous Integration` accepts `merge_group` events and runs `e2e-full` for
them, but `End-to-End Tests (full merge gate)` is not currently a required
context. Therefore the repository is prepared to test a merge-queue SHA, but
the workflow alone does not prove that every merge passed that lane.
## Pull request and merge matrix
Budgets below are job `timeout-minutes`, not typical runtimes. “Report-only”
means the result is visible and actionable but is not in the live required
context list.
| Event | Validation | Budget | Merge status | Reproduction |
|---|---|---:|---|---|
| PR, non-doc change | `Quick Checks` | 10 min | Required | `make pre-commit` (broader local umbrella) |
| PR, non-doc change | `Test and Lint` | 90 min | Required | `cargo nextest run --profile ci --all --exclude e2e_test` |
| PR, non-doc change | `Typos` | 10 min | Report-only | `typos` |
| PR, non-doc change | `ILM Integration (serial)` | 90 min | Report-only | Use the exact command in `.github/workflows/ci.yml` |
| PR, non-doc change | rio-v2 / swift / sftp test-and-lint variants | 90 min each | Report-only | `cargo nextest run` with the workflow's feature set |
| PR, non-doc change | `Build RustFS Debug Binary` | 30 min | Report-only; prerequisite for black-box lanes | `cargo build -p rustfs --bins` |
| PR, non-doc change | `io_uring Integration (real)` | 30 min | Report-only | `cargo test -p rustfs-ecstore --lib uring_ -- --test-threads=1 --nocapture` |
| PR, non-doc change | `End-to-End Tests` (`e2e-smoke` plus `s3s-e2e`) | 30 min | Report-only | `cargo nextest run --profile e2e-smoke -p e2e_test`; then `./scripts/e2e-run.sh ./target/debug/rustfs <data-dir>` |
| PR, non-doc change | `S3 Implemented Tests` | 60 min | Report-only | Build `rustfs`, then run `scripts/s3-tests/run.sh` with `DEPLOY_MODE=binary`, `TEST_MODE=single`, and `MAXFAIL=0` |
| PR, non-doc change | `S3 Lifecycle Behavior Tests` | 30 min | Report-only | Use the accelerated scanner environment in `.github/workflows/ci.yml` with `scripts/s3-tests/run.sh` |
| PR touching dependency or workflow inputs | Cargo Deny / Workflow Pin Report / Dependency Review | 20 / 5 / 30 min | Report-only | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
| PR touching architecture rules or architecture docs | `Architecture Migration Rules` | 10 min | Report-only | `scripts/check_architecture_migration_rules.sh` |
| PR touching Nix or workspace manifests | `Nix Build & Check` | 60 min | Report-only | `nix flake check` |
| PR limited to main-CI-excluded paths | companion `Quick Checks` and `Test and Lint` | 10 min each | Required | `git diff --check`; `make doc-paths-check` when documentation paths changed |
| `merge_group` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Standard required contexts only; `e2e-full` report-only | `cargo nextest run --profile e2e-full -p e2e_test` |
| Push to `main` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Post-merge detection | Same as `merge_group` |
| PR touching fuzz inputs or harness paths | Build plus five 60-second fuzz smoke targets | 60 min build; 30 min per target | Report-only | `MAX_TOTAL_TIME=60 ./scripts/fuzz/run.sh` |
| PR touching selected ecstore disk/format paths | `Rename Safety` on Windows | 60 min | Report-only | Run the four `cargo test -p rustfs-ecstore --lib <filter>` commands in `windows-filesystem.yml` on Windows |
The authoritative e2e filters live in `.config/nextest.toml`; extend a profile
instead of adding a second ad-hoc selector. Before a profile runs,
`scripts/check_test_wiring.py` compares its exact membership to the committed
digest so a silent test drop fails closed.
## Scheduled and manual validation
Scheduled lanes are independent fault domains. They do not block a pull
request, but their workflow-local gate can fail the run and scheduled failures
are routed to the shared failure-issue action. The scheduled-validation
watchdog and freshness workflow separately detect incomplete runs and missing
schedules.
| Cadence (UTC unless noted) | Workflow / validation | Budget | Verdict and artifacts | Reproduction |
|---|---|---:|---|---|
| Daily 02:17 | Fuzz: five nightly corpus targets | 60 min build; 60 min per target | Gate; corpus/crash artifacts, scheduled failure alert | `MAX_TOTAL_TIME=<seconds> ./scripts/fuzz/run.sh` |
| Daily 03:17 | MinIO interop (EC + SSE read parity) | 40 min | Gate; scheduled failure alert | Dispatch `minio-interop.yml` or follow its pinned Docker fixture steps |
| Daily 04:29 | Replication / cluster-fault / protocol e2e | 45 / 90 / 90 min | Three independent gates; JUnit, membership, and server logs | `cargo nextest run --profile e2e-repl-nightly -p e2e_test`; `--profile e2e-nightly`; `-j 1 --profile e2e-protocols` |
| Daily 06:31 | Warp performance A/B | 180 min | Regression budget gate; A/B summaries and server logs | `bash scripts/run_hotpath_warp_abba.sh --help` |
| Daily 00:07 Asia/Shanghai (16:07 UTC previous day) | Nightly GNU build and Vault lanes | 150 / 90 / 60 min | Build, live Vault, and HA failover gates | Use the commands and pinned Vault images in `nightly-gnu.yml` |
| Daily 03:23 | Security Audit | 20 / 5 min, plus 30 min on PR dependency review | Cargo Deny and workflow-pin gates; scheduled failure alert | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
| Daily 23:47 | Scheduled Validation Freshness | 10 min | Fails when a critical schedule was never created or is stale | Dispatch `scheduled-validation-freshness.yml` |
| Sunday 00:11 | Full `Continuous Integration` matrix | Per-job budgets above | Weekly variant coverage, including dormant rio-v2 binary/e2e lanes | Dispatch `ci.yml` |
| Sunday 01:13 | Seven-platform build matrix | 150 min per platform | Build/package integrity; scheduled failure alert | Dispatch `build.yml` with an exact platform set |
| Sunday 02:19 | Ceph s3-tests full sweep: single and real four-node, four shards each | 180 min per shard | Compatibility gate; report, JUnit, exact node IDs, and server logs | `scripts/s3-tests/run.sh` against an existing single or distributed target |
| Sunday 06:41 | Mint | 120 min | **Report-only by design**; per-suite PASS/FAIL/NA and raw `log.json` | Reproduce the pinned Docker sequence in `mint.yml` or dispatch it |
| Sunday 07:43 | Workspace line coverage | 120 min | Report-only trend; lcov and JSON retained 90 days | `make coverage` |
| Monthly, day 1 06:37 | Runner Hygiene | 15 min | Validates runner ephemerality; scheduled failure alert | Dispatch `runner-hygiene.yml` |
Manual `workflow_dispatch` exists for the scheduled workflows above. Manual
runs are debugging evidence and intentionally do not open scheduled-failure
issues. A manual performance run may explicitly allow a known regression; that
override must not be treated as an ordinary passing baseline.
## Release validation
Release validation is post-merge and tag-driven; it does not substitute for a
pull-request gate.
| Event | Validation | Budget | Result |
|---|---|---:|---|
| Push to `main` or weekly schedule | `Build and Release` platform matrix | 150 min per platform | Build artifacts for all selected targets; no release publication on a main push |
| Valid release or preview tag | `Build and Release` plus asset checks | 150 min per platform | Draft release, checksummed assets, and publish step |
| Successful non-preview release-tag build | Docker image build and image scan | 60 min build; 30 min scan | Multi-architecture images plus vulnerability report |
| Successful release-tag build | DEB/RPM packaging | 30 min per architecture | Packages and checksum files uploaded to the release |
| Successful non-preview release-tag build | Helm template test and package | 30 min build; 30 min publish | Versioned chart and repository index |
Use an exact preview tag for end-to-end release rehearsal. Manual dispatches
are backfill/debug paths and do not prove the automatic `workflow_run` chain.
## Evidence requirements
A green check is useful only when it proves the intended behavior ran:
- Record the exact commit SHA and run URL.
- Separate product failure from runner prerequisites, service readiness, and
cancellation. Repair the precondition, then rerun the exact workload.
- Preserve membership manifests, JUnit, raw compatibility logs, seeds, and
server logs where the workflow provides them.
- For a bug fix or a new fault checker, provide sensitivity evidence: the old
behavior or an intentional mutation must fail the new oracle, and the fixed
behavior must pass it.
- Never promote a report-only lane to required from one green run. Require at
least 14 days and 30 representative pull requests with at least 99% complete
execution, then update the ruleset and this table together.
## Change checklist
Update this file in the same pull request when any of these change:
- workflow triggers, job names, timeouts, or nextest profile ownership;
- required status contexts or strict/merge-queue policy;
- scheduled cadence, alert routing, artifact contract, or local reproduction;
- report-only versus gating semantics.
Do not copy per-module test counts here. Update
[e2e-suite-inventory.md](e2e-suite-inventory.md) and its enforced membership
digest instead.
+2 -2
View File
@@ -336,13 +336,13 @@ opentelemetry = { workspace = true }
tracing-opentelemetry = { workspace = true }
# Data structures
hashbrown = { workspace = true, features = ["serde", "rayon"] }
mimalloc = { workspace = true }
rustfs-mimalloc = { workspace = true }
[target.'cfg(target_os = "linux")'.dependencies]
libsystemd.workspace = true
[target.'cfg(not(target_os = "windows"))'.dependencies]
libmimalloc-sys.workspace = true
rustfs-mimalloc-sys.workspace = true
[dev-dependencies]
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
+1 -7
View File
@@ -369,14 +369,8 @@ pub fn allocator_reclaim_controller_snapshot(ctx: &CancellationToken) -> Allocat
}
#[cfg(not(target_os = "windows"))]
#[allow(unsafe_code)]
fn collect_allocator_memory(force: bool) -> Result<(), String> {
// SAFETY: `mi_collect` is provided by the active global allocator backend
// on this target family. It is explicitly intended to reclaim retained
// pages/segments and does not require additional invariants from the caller.
unsafe {
libmimalloc_sys::mi_collect(force);
}
rustfs_mimalloc::MiMalloc::collect(force);
Ok(())
}
+187 -44
View File
@@ -19,19 +19,23 @@ use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use chrono::{DateTime, SecondsFormat, Utc};
use reqwest::{Client, StatusCode, Url, header};
use rustls::RootCertStore;
use rustls::pki_types::{CertificateDer, pem::PemObject as _};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use zeroize::Zeroizing;
use super::config::HeartbeatConfig;
use super::credential_store::CredentialStoreError;
use super::credential_store::{CredentialStoreError, DeviceCredential};
use super::identity::IdentityError;
use super::identity_store::StoreError;
use super::registration::CredentialValidationError;
use super::telemetry::{TelemetryDelivery, TelemetryError, TelemetryTransport, is_exact_utc_seconds};
use super::registration::{CredentialValidationError, validate_stored_credential};
const PROTOCOL_VERSION: &str = "v1";
const AGENT_VERSION: &str = concat!("rustfs-agent/", env!("CARGO_PKG_VERSION"));
const MAX_SEQUENCE: u64 = 9_007_199_254_740_991;
const MAX_RESPONSE_BYTES: usize = 64 * 1024;
#[cfg(unix)]
const FILE_MODE: u32 = 0o600;
static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
@@ -118,39 +122,133 @@ pub(crate) enum Delivery {
}
pub(crate) struct HeartbeatSender {
transport: TelemetryTransport,
endpoint: Url,
root_store: RootCertStore,
roots: Vec<CertificateDer<'static>>,
config: HeartbeatConfig,
}
impl HeartbeatSender {
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, HeartbeatError> {
let mut endpoint = Url::parse(&config.endpoint).map_err(|_| HeartbeatError::Endpoint)?;
if endpoint.scheme() != "https"
|| endpoint.cannot_be_a_base()
|| !endpoint.username().is_empty()
|| endpoint.password().is_some()
|| endpoint.query().is_some()
|| endpoint.fragment().is_some()
{
return Err(HeartbeatError::Endpoint);
}
if !endpoint.path().ends_with('/') {
endpoint.set_path(&format!("{}/", endpoint.path()));
}
let roots = CertificateDer::pem_slice_iter(&config.root_ca_pem)
.collect::<Result<Vec<_>, _>>()
.map_err(|_| HeartbeatError::RootCertificate)?;
if roots.is_empty() {
return Err(HeartbeatError::RootCertificate);
}
let mut root_store = RootCertStore::empty();
let (accepted, rejected) = root_store.add_parsable_certificates(roots.clone());
if accepted != roots.len() || rejected != 0 {
return Err(HeartbeatError::RootCertificate);
}
let schedule = config.schedule;
if schedule.cadence.is_zero() || schedule.jitter > schedule.cadence {
if schedule.cadence.is_zero()
|| schedule.timeout.is_zero()
|| schedule.timeout > Duration::from_secs(5)
|| schedule.initial_backoff.is_zero()
|| schedule.max_backoff < schedule.initial_backoff
|| schedule.max_backoff > Duration::from_secs(5 * 60)
|| schedule.jitter > schedule.cadence
{
return Err(HeartbeatError::Schedule);
}
Ok(Self {
transport: TelemetryTransport::new(config)?,
endpoint,
root_store,
roots,
config,
})
}
pub(crate) async fn send(&self, heartbeat: &PendingHeartbeat) -> Result<Delivery, HeartbeatError> {
match self.transport.post("heartbeats", heartbeat).await? {
TelemetryDelivery::Accepted { body, .. } => {
let accepted: HeartbeatResponse = serde_json::from_slice(&body).map_err(|_| HeartbeatError::Response)?;
if accepted.accepted_version != PROTOCOL_VERSION
|| accepted.capability_hints.len() > 32
|| accepted.capability_hints.iter().any(|hint| hint.len() > 32)
|| !is_exact_utc_seconds(&accepted.server_time)
{
return Err(HeartbeatError::Response);
}
Ok(Delivery::Accepted {
server_time: accepted.server_time,
})
let (cluster_uid, client) = {
let _lock = self.config.credential_store.lock().await?;
let credential = self.config.credential_store.load()?.ok_or(HeartbeatError::NotRegistered)?;
let identity = self.config.identity_store.load()?.ok_or(HeartbeatError::IdentityMissing)?;
validate_stored_credential(&credential, &identity, &self.root_store, &self.roots)?;
let now = Utc::now().timestamp();
if now < credential.not_before_unix || now >= credential.not_after_unix {
return Err(HeartbeatError::CredentialExpired);
}
TelemetryDelivery::Retry { retry_after } => Ok(Delivery::Retry { retry_after }),
TelemetryDelivery::AuthenticationStopped { status, reason } => Ok(Delivery::AuthenticationStopped { status, reason }),
TelemetryDelivery::Rejected { status, reason } => Ok(Delivery::Rejected { status, reason }),
let cluster_uid = cluster_uid(&credential)?.to_owned();
let client = self.client(&credential, &identity.to_pkcs8_pem()?)?;
(cluster_uid, client)
};
let url = self.endpoint.join(&format!("clusters/{cluster_uid}/heartbeats"))?;
let response = match client.post(url).json(heartbeat).send().await {
Ok(response) => response,
Err(error) if error.is_timeout() || error.is_connect() || error.is_request() => {
return Ok(Delivery::Retry { retry_after: None });
}
Err(error) => return Err(error.into()),
};
let status = response.status();
if status == StatusCode::TOO_MANY_REQUESTS {
return Ok(Delivery::Retry {
retry_after: retry_after(response.headers(), Utc::now(), self.config.schedule.max_backoff),
});
}
if status == StatusCode::REQUEST_TIMEOUT || status.is_server_error() {
return Ok(Delivery::Retry { retry_after: None });
}
if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) {
return Ok(Delivery::AuthenticationStopped {
status: status.as_u16(),
reason: response_reason(response).await,
});
}
if status != StatusCode::OK {
return Ok(Delivery::Rejected {
status: status.as_u16(),
reason: response_reason(response).await,
});
}
let accepted: HeartbeatResponse =
serde_json::from_slice(&bounded_body(response).await?).map_err(|_| HeartbeatError::Response)?;
if accepted.accepted_version != PROTOCOL_VERSION
|| accepted.capability_hints.len() > 32
|| accepted.capability_hints.iter().any(|hint| hint.len() > 32)
|| !is_exact_utc_seconds(&accepted.server_time)
{
return Err(HeartbeatError::Response);
}
Ok(Delivery::Accepted {
server_time: accepted.server_time,
})
}
fn client(&self, credential: &DeviceCredential, key: &Zeroizing<String>) -> Result<Client, HeartbeatError> {
let mut pem = Zeroizing::new(Vec::with_capacity(credential.certificate_chain.len() + key.len() + 1));
pem.extend_from_slice(credential.certificate_chain.as_bytes());
pem.push(b'\n');
pem.extend_from_slice(key.as_bytes());
let identity = reqwest::Identity::from_pem(&pem).map_err(|_| HeartbeatError::IdentityCertificate)?;
let roots = self
.roots
.iter()
.map(|root| reqwest::Certificate::from_der(root.as_ref()))
.collect::<Result<Vec<_>, _>>()?;
Client::builder()
.https_only(true)
.redirect(reqwest::redirect::Policy::none())
.timeout(self.config.schedule.timeout)
.tls_certs_only(roots)
.identity(identity)
.build()
.map_err(Into::into)
}
}
@@ -280,6 +378,73 @@ impl HeartbeatStateStore {
}
}
fn cluster_uid(credential: &DeviceCredential) -> Result<&str, HeartbeatError> {
let mut parts = credential.name.split('/');
let valid = parts.next() == Some("organizations");
let organization_uid = parts.next();
let valid = valid && parts.next() == Some("clusters");
let cluster_uid = parts.next();
let valid = valid && parts.next() == Some("clusterDevices");
let device_uid = parts.next();
if !valid
|| organization_uid.is_none_or(str::is_empty)
|| cluster_uid.is_none_or(str::is_empty)
|| device_uid != Some(credential.uid.as_str())
|| parts.next().is_some()
{
return Err(HeartbeatError::CredentialName);
}
cluster_uid.ok_or(HeartbeatError::CredentialName)
}
fn retry_after(headers: &header::HeaderMap, now: DateTime<Utc>, maximum: Duration) -> Option<Duration> {
let value = headers.get(header::RETRY_AFTER)?.to_str().ok()?;
let delay = value.parse::<u64>().ok().map(Duration::from_secs).or_else(|| {
DateTime::parse_from_rfc2822(value)
.ok()
.and_then(|at| (at.with_timezone(&Utc) - now).to_std().ok())
})?;
Some(delay.min(maximum))
}
fn is_exact_utc_seconds(value: &str) -> bool {
DateTime::parse_from_rfc3339(value).is_ok_and(|time| {
time.offset().local_minus_utc() == 0
&& value.ends_with('Z')
&& time.with_timezone(&Utc).to_rfc3339_opts(SecondsFormat::Secs, true) == value
})
}
async fn response_reason(response: reqwest::Response) -> Option<String> {
#[derive(Deserialize)]
struct Envelope {
#[serde(default)]
details: Vec<Detail>,
}
#[derive(Deserialize)]
struct Detail {
#[serde(default)]
reason: String,
}
serde_json::from_slice::<Envelope>(&bounded_body(response).await.ok()?)
.ok()?
.details
.into_iter()
.find_map(|detail| (!detail.reason.is_empty()).then_some(detail.reason))
}
async fn bounded_body(mut response: reqwest::Response) -> Result<Vec<u8>, HeartbeatError> {
let mut body = Vec::new();
while let Some(chunk) = response.chunk().await? {
if body.len().saturating_add(chunk.len()) > MAX_RESPONSE_BYTES {
return Err(HeartbeatError::ResponseTooLarge);
}
body.extend_from_slice(&chunk);
}
Ok(body)
}
fn parent(path: &Path) -> Result<&Path, HeartbeatError> {
path.parent()
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state path has no parent")))
@@ -418,25 +583,3 @@ pub enum HeartbeatError {
#[error(transparent)]
CredentialValidation(#[from] CredentialValidationError),
}
impl From<TelemetryError> for HeartbeatError {
fn from(error: TelemetryError) -> Self {
match error {
TelemetryError::Endpoint => Self::Endpoint,
TelemetryError::RootCertificate => Self::RootCertificate,
TelemetryError::Schedule => Self::Schedule,
TelemetryError::NotRegistered => Self::NotRegistered,
TelemetryError::IdentityMissing => Self::IdentityMissing,
TelemetryError::IdentityCertificate => Self::IdentityCertificate,
TelemetryError::CredentialName => Self::CredentialName,
TelemetryError::CredentialExpired => Self::CredentialExpired,
TelemetryError::ResponseTooLarge => Self::ResponseTooLarge,
TelemetryError::Url(error) => Self::Url(error),
TelemetryError::Transport(error) => Self::Transport(error),
TelemetryError::Identity(error) => Self::Identity(error),
TelemetryError::IdentityStore(error) => Self::IdentityStore(error),
TelemetryError::CredentialStore(error) => Self::CredentialStore(error),
TelemetryError::CredentialValidation(error) => Self::CredentialValidation(error),
}
}
}
-594
View File
@@ -1,594 +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 std::collections::BTreeSet;
use std::fs;
use std::io::{self, Write as _};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use serde::{Deserialize, Serialize};
use sha2::{Digest as _, Sha256};
use uuid::Uuid;
use super::config::HeartbeatConfig;
use super::telemetry::{TelemetryDelivery, TelemetryError, TelemetryTransport, is_exact_utc_seconds};
const PROTOCOL_VERSION: &str = "v1";
const RUSTFS_VERSION: &str = concat!(
env!("CARGO_PKG_VERSION_MAJOR"),
".",
env!("CARGO_PKG_VERSION_MINOR"),
".",
env!("CARGO_PKG_VERSION_PATCH")
);
const HASH_PREFIX: &[u8] = b"rustfs-connect/agent/v1/inventory-snapshot\n";
const MAX_SEQUENCE: u64 = 9_007_199_254_740_991;
const MAX_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
#[cfg(unix)]
const FILE_MODE: u32 = 0o600;
static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct InventorySchedule {
pub cadence: Duration,
pub jitter: Duration,
}
impl Default for InventorySchedule {
fn default() -> Self {
Self {
cadence: Duration::from_secs(6 * 60 * 60),
jitter: Duration::from_secs(30 * 60),
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum InventoryStatus {
Starting,
Unchanged { content_hash: String },
Online { content_hash: String, received_at: String },
BackingOff { delay: Duration },
AuthenticationStopped { status: u16, reason: Option<String> },
Failed { reason: String },
Stopped,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub enum InventoryFlag {
#[serde(rename = "capacity.critical")]
CapacityCritical,
#[serde(rename = "capacity.warning")]
CapacityWarning,
#[serde(rename = "clock.skew")]
ClockSkew,
#[serde(rename = "cluster.degraded")]
ClusterDegraded,
#[serde(rename = "cluster.healing")]
ClusterHealing,
#[serde(rename = "cluster.readonly")]
ClusterReadonly,
#[serde(rename = "drive.offline")]
DriveOffline,
#[serde(rename = "node.offline")]
NodeOffline,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum OperatingSystemFamily {
Linux,
Darwin,
Windows,
Freebsd,
Other,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct InventoryOsVersion {
family: OperatingSystemFamily,
major: u16,
minor: u16,
}
impl InventoryOsVersion {
pub fn new(family: OperatingSystemFamily, major: u16, minor: u16) -> Result<Self, InventoryError> {
if major > 9999 || minor > 9999 {
return Err(InventoryError::OsVersion);
}
Ok(Self { family, major, minor })
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct InventorySnapshot {
rustfs_version: String,
os_version: Option<InventoryOsVersion>,
node_count: u16,
drive_count: u32,
capacity_total_bytes: u64,
capacity_used_bytes: u64,
coarse_flags: Vec<InventoryFlag>,
}
impl InventorySnapshot {
pub fn current(
node_count: usize,
drive_count: usize,
capacity_total_bytes: u64,
capacity_free_bytes: u64,
coarse_flags: impl IntoIterator<Item = InventoryFlag>,
) -> Result<Self, InventoryError> {
let capacity_used_bytes = capacity_total_bytes
.checked_sub(capacity_free_bytes)
.ok_or(InventoryError::Capacity)?;
Self::new(
RUSTFS_VERSION,
None,
node_count,
drive_count,
capacity_total_bytes,
capacity_used_bytes,
coarse_flags,
)
}
pub fn new(
rustfs_version: impl Into<String>,
os_version: Option<InventoryOsVersion>,
node_count: usize,
drive_count: usize,
capacity_total_bytes: u64,
capacity_used_bytes: u64,
coarse_flags: impl IntoIterator<Item = InventoryFlag>,
) -> Result<Self, InventoryError> {
let snapshot = Self {
rustfs_version: rustfs_version.into(),
os_version,
node_count: u16::try_from(node_count).map_err(|_| InventoryError::NodeCount)?,
drive_count: u32::try_from(drive_count).map_err(|_| InventoryError::DriveCount)?,
capacity_total_bytes,
capacity_used_bytes,
coarse_flags: coarse_flags.into_iter().collect::<BTreeSet<_>>().into_iter().collect(),
};
snapshot.validate()?;
Ok(snapshot)
}
pub fn content_hash(&self) -> Result<String, InventoryError> {
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct Canonical<'a> {
capacity_total_bytes: u64,
capacity_used_bytes: u64,
coarse_flags: &'a [InventoryFlag],
drive_count: u32,
node_count: u16,
os_version: Option<InventoryOsVersion>,
rustfs_version: &'a str,
}
let canonical = serde_json::to_vec(&Canonical {
capacity_total_bytes: self.capacity_total_bytes,
capacity_used_bytes: self.capacity_used_bytes,
coarse_flags: &self.coarse_flags,
drive_count: self.drive_count,
node_count: self.node_count,
os_version: self.os_version,
rustfs_version: &self.rustfs_version,
})?;
let mut digest = Sha256::new();
digest.update(HASH_PREFIX);
digest.update(canonical);
Ok(hex_simd::encode_to_string(digest.finalize(), hex_simd::AsciiCase::Lower))
}
fn validate(&self) -> Result<(), InventoryError> {
if !valid_version(&self.rustfs_version) {
return Err(InventoryError::RustfsVersion);
}
if self.node_count == 0 || self.node_count > 4096 {
return Err(InventoryError::NodeCount);
}
if self.drive_count > 1_048_576 {
return Err(InventoryError::DriveCount);
}
if self.capacity_total_bytes > MAX_SAFE_INTEGER || self.capacity_used_bytes > self.capacity_total_bytes {
return Err(InventoryError::Capacity);
}
Ok(())
}
}
fn valid_version(version: &str) -> bool {
let components = version.split('.').collect::<Vec<_>>();
components.len() == 3
&& components.iter().all(|component| {
!component.is_empty()
&& component.len() <= 4
&& (component == &"0" || !component.starts_with('0'))
&& component.parse::<u16>().is_ok_and(|value| value <= 9999)
})
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub(crate) struct PendingInventory {
protocol_version: String,
request_id: String,
sequence: u64,
#[serde(flatten)]
snapshot: InventorySnapshot,
}
impl PendingInventory {
fn new(snapshot: InventorySnapshot, sequence: u64) -> Self {
Self {
protocol_version: PROTOCOL_VERSION.to_owned(),
request_id: Uuid::new_v4().to_string(),
sequence,
snapshot,
}
}
fn is_valid(&self) -> bool {
self.protocol_version == PROTOCOL_VERSION
&& self.sequence <= MAX_SEQUENCE
&& self.snapshot.validate().is_ok()
&& Uuid::parse_str(&self.request_id)
.is_ok_and(|request_id| request_id.get_version_num() == 4 && request_id.to_string() == self.request_id)
}
fn content_hash(&self) -> Result<String, InventoryError> {
self.snapshot.content_hash()
}
}
pub(crate) enum InventoryDelivery {
Accepted { content_hash: String, received_at: String },
Retry { retry_after: Option<Duration> },
AuthenticationStopped { status: u16, reason: Option<String> },
Rejected { status: u16, reason: Option<String> },
}
pub(crate) struct InventorySender {
transport: TelemetryTransport,
}
impl InventorySender {
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, InventoryError> {
Ok(Self {
transport: TelemetryTransport::new(config)?,
})
}
pub(crate) async fn send(&self, inventory: &PendingInventory) -> Result<InventoryDelivery, InventoryError> {
match self.transport.post("inventorySnapshots", inventory).await? {
TelemetryDelivery::Accepted { cluster_name, body } => {
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct InventoryResponse {
name: String,
uid: String,
content_hash: String,
received_at: String,
}
let accepted: InventoryResponse = serde_json::from_slice(&body).map_err(|_| InventoryError::Response)?;
let uid = Uuid::parse_str(&accepted.uid).map_err(|_| InventoryError::Response)?;
let content_hash = inventory.content_hash()?;
if uid.get_version_num() != 7
|| uid.to_string() != accepted.uid
|| accepted.name != format!("{cluster_name}/inventorySnapshots/{}", accepted.uid)
|| accepted.content_hash != content_hash
|| !is_exact_utc_seconds(&accepted.received_at)
{
return Err(InventoryError::Response);
}
Ok(InventoryDelivery::Accepted {
content_hash,
received_at: accepted.received_at,
})
}
TelemetryDelivery::Retry { retry_after } => Ok(InventoryDelivery::Retry { retry_after }),
TelemetryDelivery::AuthenticationStopped { status, reason } => {
Ok(InventoryDelivery::AuthenticationStopped { status, reason })
}
TelemetryDelivery::Rejected { status, reason } => Ok(InventoryDelivery::Rejected { status, reason }),
}
}
}
#[derive(Clone)]
pub(crate) struct InventoryStateStore {
path: PathBuf,
}
#[derive(Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
struct InventoryState {
next_sequence: u64,
pending: Option<PendingInventory>,
last_accepted_content_hash: Option<String>,
}
impl InventoryStateStore {
pub(crate) fn from_heartbeat_path(path: &Path) -> Result<Self, InventoryError> {
let root = path.parent().and_then(Path::parent).ok_or(InventoryError::StatePath)?;
Ok(Self {
path: root.join("inventory/state.json"),
})
}
pub(crate) fn try_runtime_lock(&self) -> Result<fs::File, InventoryError> {
let directory = parent(&self.path)?;
fs::create_dir_all(directory).map_err(|source| state_io(directory, source))?;
let name = filename(&self.path)?;
let path = directory.join(format!(".{name}.lock"));
let mut options = fs::OpenOptions::new();
options.create(true).truncate(false).read(true).write(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(FILE_MODE);
}
let lock = options.open(&path).map_err(|source| state_io(&path, source))?;
check_mode(&path)?;
lock.try_lock().map_err(|_| InventoryError::AlreadyRunning)?;
Ok(lock)
}
pub(crate) async fn pending(&self) -> Result<Option<PendingInventory>, InventoryError> {
let store = self.clone();
tokio::task::spawn_blocking(move || {
let state = store.read()?;
if state.pending.is_none() && state.next_sequence > MAX_SEQUENCE {
return Err(InventoryError::SequenceExhausted);
}
Ok(state.pending)
})
.await
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
}
pub(crate) async fn prepare(&self, snapshot: InventorySnapshot) -> Result<Option<PendingInventory>, InventoryError> {
let store = self.clone();
tokio::task::spawn_blocking(move || store.prepare_sync(snapshot))
.await
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
}
pub(crate) async fn mark_accepted(&self, accepted: &PendingInventory) -> Result<(), InventoryError> {
let store = self.clone();
let accepted = accepted.clone();
tokio::task::spawn_blocking(move || store.mark_accepted_sync(&accepted))
.await
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
}
fn prepare_sync(&self, snapshot: InventorySnapshot) -> Result<Option<PendingInventory>, InventoryError> {
let mut state = self.read()?;
if state.pending.is_some() {
return Ok(state.pending);
}
let content_hash = snapshot.content_hash()?;
if state.last_accepted_content_hash.as_deref() == Some(&content_hash) {
return Ok(None);
}
if state.next_sequence > MAX_SEQUENCE {
return Err(InventoryError::SequenceExhausted);
}
let pending = PendingInventory::new(snapshot, state.next_sequence);
state.pending = Some(pending.clone());
self.write(&state)?;
Ok(Some(pending))
}
fn mark_accepted_sync(&self, accepted: &PendingInventory) -> Result<(), InventoryError> {
let mut state = self.read()?;
if state.pending.as_ref() != Some(accepted) {
return Err(InventoryError::StateConflict);
}
state.next_sequence = accepted.sequence.checked_add(1).ok_or(InventoryError::SequenceExhausted)?;
state.last_accepted_content_hash = Some(accepted.content_hash()?);
state.pending = None;
self.write(&state)
}
fn read(&self) -> Result<InventoryState, InventoryError> {
let bytes = match fs::read(&self.path) {
Ok(bytes) => bytes,
Err(source) if source.kind() == io::ErrorKind::NotFound => return Ok(InventoryState::default()),
Err(source) => return Err(state_io(&self.path, source)),
};
check_mode(&self.path)?;
let state: InventoryState = serde_json::from_slice(&bytes).map_err(|source| InventoryError::StateInvalid {
path: self.path.clone(),
source,
})?;
let last_hash_valid = state.last_accepted_content_hash.as_deref().is_none_or(valid_content_hash);
let pending_valid = state.pending.as_ref().is_none_or(|pending| {
pending.sequence == state.next_sequence
&& pending.is_valid()
&& pending
.content_hash()
.is_ok_and(|hash| state.last_accepted_content_hash.as_deref() != Some(&hash))
});
if state.next_sequence > MAX_SEQUENCE + 1 || !last_hash_valid || !pending_valid {
return Err(InventoryError::StateCorrupt { path: self.path.clone() });
}
Ok(state)
}
fn write(&self, state: &InventoryState) -> Result<(), InventoryError> {
let bytes = serde_json::to_vec(state).map_err(|source| InventoryError::StateInvalid {
path: self.path.clone(),
source,
})?;
let directory = parent(&self.path)?;
fs::create_dir_all(directory).map_err(|source| state_io(directory, source))?;
let temp = stage(directory, &self.path, &bytes)?;
let result = fs::rename(&temp, &self.path)
.map_err(|source| state_io(&self.path, source))
.and_then(|()| fsync_dir(directory).map_err(|source| state_io(directory, source)));
if result.is_err() {
let _ = fs::remove_file(temp);
}
result
}
}
fn valid_content_hash(value: &str) -> bool {
value.len() == 64
&& value
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
}
fn parent(path: &Path) -> Result<&Path, InventoryError> {
path.parent()
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state path has no parent")))
}
fn filename(path: &Path) -> Result<&str, InventoryError> {
path.file_name()
.and_then(|name| name.to_str())
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state filename is invalid")))
}
fn stage(directory: &Path, destination: &Path, bytes: &[u8]) -> Result<PathBuf, InventoryError> {
let name = filename(destination)?;
loop {
let path = directory.join(format!(
".{name}.{}.{}.tmp",
std::process::id(),
STAGING_SEQUENCE.fetch_add(1, Ordering::Relaxed)
));
let mut options = fs::OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(FILE_MODE);
}
let mut file = match options.open(&path) {
Ok(file) => file,
Err(source) if source.kind() == io::ErrorKind::AlreadyExists => continue,
Err(source) => return Err(state_io(&path, source)),
};
if let Err(source) = file.write_all(bytes).and_then(|()| file.sync_all()) {
let _ = fs::remove_file(&path);
return Err(state_io(&path, source));
}
return Ok(path);
}
}
fn state_io(path: &Path, source: io::Error) -> InventoryError {
InventoryError::StateIo {
path: path.to_path_buf(),
source,
}
}
#[cfg(unix)]
fn check_mode(path: &Path) -> Result<(), InventoryError> {
use std::os::unix::fs::PermissionsExt as _;
let mode = fs::metadata(path)
.map_err(|source| state_io(path, source))?
.permissions()
.mode()
& 0o7777;
if mode != FILE_MODE {
return Err(InventoryError::StatePermissions {
path: path.to_path_buf(),
mode,
expected: FILE_MODE,
});
}
Ok(())
}
#[cfg(not(unix))]
fn check_mode(_path: &Path) -> Result<(), InventoryError> {
Ok(())
}
fn fsync_dir(directory: &Path) -> io::Result<()> {
#[cfg(unix)]
fs::File::open(directory)?.sync_all()?;
#[cfg(not(unix))]
let _ = directory;
Ok(())
}
#[derive(Debug, thiserror::Error)]
pub enum InventoryError {
#[error("the RustFS inventory version is outside protocol bounds")]
RustfsVersion,
#[error("the RustFS inventory operating-system version is outside protocol bounds")]
OsVersion,
#[error("the RustFS inventory node count is outside protocol bounds")]
NodeCount,
#[error("the RustFS inventory drive count is outside protocol bounds")]
DriveCount,
#[error("the RustFS inventory capacity is outside protocol bounds")]
Capacity,
#[error("the RustFS inventory snapshot is incomplete: observed {observed} of {expected} configured drives")]
SnapshotIncomplete { expected: usize, observed: usize },
#[error("the Connect inventory schedule is invalid")]
Schedule,
#[error("the Connect inventory sequence is exhausted")]
SequenceExhausted,
#[error("a Connect inventory runtime already owns this state")]
AlreadyRunning,
#[error("the persisted Connect inventory changed while delivery was in flight")]
StateConflict,
#[error("the Connect inventory state path is invalid")]
StatePath,
#[error("Connect inventory state I/O failed at {path}: {source}")]
StateIo {
path: PathBuf,
#[source]
source: io::Error,
},
#[error("Connect inventory state at {path} is invalid: {source}")]
StateInvalid {
path: PathBuf,
#[source]
source: serde_json::Error,
},
#[error("Connect inventory state at {path} violates the protocol invariants")]
StateCorrupt { path: PathBuf },
#[cfg(unix)]
#[error("Connect inventory state at {path} has mode {mode:o}, expected {expected:o}")]
StatePermissions { path: PathBuf, mode: u32, expected: u32 },
#[error("Connect returned an invalid inventory response")]
Response,
#[error(transparent)]
Json(#[from] serde_json::Error),
#[error("Connect inventory delivery failed: {0}")]
Telemetry(String),
}
impl From<TelemetryError> for InventoryError {
fn from(error: TelemetryError) -> Self {
Self::Telemetry(error.to_string())
}
}
+1 -7
View File
@@ -31,11 +31,9 @@ pub mod credential_store;
pub mod heartbeat;
pub mod identity;
pub mod identity_store;
pub mod inventory;
pub mod offline;
pub mod registration;
pub mod runtime;
mod telemetry;
pub use client::{ClientError, ConnectClient, ConnectConfig};
pub use config::{HeartbeatConfig, HeartbeatConfigError, HeartbeatSchedule};
@@ -43,10 +41,6 @@ pub use credential_store::{CredentialStore, DeviceCredential};
pub use heartbeat::{CoarseNodeSummary, HeartbeatError, HeartbeatStatus};
pub use identity::{DeviceIdentity, IdentityError, RegistrationProof, RegistrationTranscript};
pub use identity_store::{IdentityStore, StoreError};
pub use inventory::{
InventoryError, InventoryFlag, InventoryOsVersion, InventorySchedule, InventorySnapshot, InventoryStatus,
OperatingSystemFamily,
};
pub use offline::{EnrollmentError, OfflineEnrollment, OfflineKeyStore, VerifiedChallenge};
pub use registration::{RegistrationToken, TokenError};
pub use runtime::{HeartbeatRuntime, InventoryRuntime, spawn_heartbeat_runtime, spawn_inventory_runtime};
pub use runtime::{HeartbeatRuntime, spawn_heartbeat_runtime};
-164
View File
@@ -23,16 +23,11 @@ use tokio_util::sync::CancellationToken;
use super::config::HeartbeatConfig;
use super::heartbeat::{CoarseNodeSummary, Delivery, HeartbeatError, HeartbeatSender, HeartbeatStateStore, HeartbeatStatus};
use super::inventory::{
InventoryDelivery, InventoryError, InventorySchedule, InventorySender, InventorySnapshot, InventoryStateStore,
InventoryStatus,
};
pub struct HeartbeatRuntime {
shutdown: CancellationToken,
status: watch::Receiver<HeartbeatStatus>,
task: Option<JoinHandle<()>>,
inventory: Option<InventoryRuntime>,
}
impl HeartbeatRuntime {
@@ -40,19 +35,11 @@ impl HeartbeatRuntime {
self.status.clone()
}
pub(crate) fn with_inventory(mut self, inventory: Option<InventoryRuntime>) -> Self {
self.inventory = inventory;
self
}
pub async fn shutdown(mut self) {
self.shutdown.cancel();
if let Some(task) = self.task.take() {
let _ = task.await;
}
if let Some(inventory) = self.inventory.take() {
inventory.shutdown().await;
}
}
}
@@ -62,31 +49,6 @@ impl Drop for HeartbeatRuntime {
}
}
pub struct InventoryRuntime {
shutdown: CancellationToken,
status: watch::Receiver<InventoryStatus>,
task: Option<JoinHandle<()>>,
}
impl InventoryRuntime {
pub fn status(&self) -> watch::Receiver<InventoryStatus> {
self.status.clone()
}
pub async fn shutdown(mut self) {
self.shutdown.cancel();
if let Some(task) = self.task.take() {
let _ = task.await;
}
}
}
impl Drop for InventoryRuntime {
fn drop(&mut self) {
self.shutdown.cancel();
}
}
pub fn spawn_heartbeat_runtime<F>(
config: Option<HeartbeatConfig>,
parent_shutdown: &CancellationToken,
@@ -160,126 +122,6 @@ where
shutdown,
status: status_rx,
task: Some(task),
inventory: None,
}))
}
pub fn spawn_inventory_runtime<F, Fut>(
config: Option<HeartbeatConfig>,
schedule: InventorySchedule,
parent_shutdown: &CancellationToken,
sample: F,
) -> Result<Option<InventoryRuntime>, InventoryError>
where
F: Fn() -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<InventorySnapshot, InventoryError>> + Send + 'static,
{
let Some(config) = config else {
return Ok(None);
};
if schedule.cadence.is_zero() || schedule.jitter > schedule.cadence {
return Err(InventoryError::Schedule);
}
let retry_schedule = config.schedule;
let store = InventoryStateStore::from_heartbeat_path(&config.state_path)?;
let lock = store.try_runtime_lock()?;
let sender = InventorySender::new(config)?;
let shutdown = parent_shutdown.child_token();
let task_shutdown = shutdown.clone();
let (status_tx, status_rx) = watch::channel(InventoryStatus::Starting);
let task = tokio::spawn(async move {
let _lock = lock;
let mut backoff = retry_schedule.initial_backoff;
loop {
if task_shutdown.is_cancelled() {
break;
}
let pending = match store.pending().await {
Ok(Some(pending)) => pending,
Ok(None) => {
let snapshot = match cancellable(&task_shutdown, sample()).await {
Some(Ok(snapshot)) => snapshot,
Some(Err(InventoryError::SnapshotIncomplete { .. })) => {
let delay = backoff;
backoff = backoff.saturating_mul(2).min(retry_schedule.max_backoff);
let _ = status_tx.send(InventoryStatus::BackingOff { delay });
if sleep_or_cancel(&task_shutdown, delay).await {
break;
}
continue;
}
Some(Err(error)) => return failed_inventory(&status_tx, error),
None => break,
};
let content_hash = match snapshot.content_hash() {
Ok(content_hash) => content_hash,
Err(error) => return failed_inventory(&status_tx, error),
};
match store.prepare(snapshot).await {
Ok(Some(pending)) => pending,
Ok(None) => {
backoff = retry_schedule.initial_backoff;
let _ = status_tx.send(InventoryStatus::Unchanged { content_hash });
if sleep_or_cancel(&task_shutdown, schedule.cadence.saturating_add(jitter(schedule.jitter))).await {
break;
}
continue;
}
Err(error) => return failed_inventory(&status_tx, error),
}
}
Err(error) => return failed_inventory(&status_tx, error),
};
let delivery = match cancellable(&task_shutdown, sender.send(&pending)).await {
Some(Ok(delivery)) => delivery,
Some(Err(error)) => return failed_inventory(&status_tx, error),
None => break,
};
let delay = match delivery {
InventoryDelivery::Accepted {
content_hash,
received_at,
} => {
if let Err(error) = store.mark_accepted(&pending).await {
return failed_inventory(&status_tx, error);
}
backoff = retry_schedule.initial_backoff;
let _ = status_tx.send(InventoryStatus::Online {
content_hash,
received_at,
});
schedule.cadence.saturating_add(jitter(schedule.jitter))
}
InventoryDelivery::Retry { retry_after } => {
let delay = retry_after
.unwrap_or(backoff)
.clamp(retry_schedule.initial_backoff, retry_schedule.max_backoff);
backoff = backoff.saturating_mul(2).min(retry_schedule.max_backoff);
let _ = status_tx.send(InventoryStatus::BackingOff { delay });
delay
}
InventoryDelivery::AuthenticationStopped { status, reason } => {
let _ = status_tx.send(InventoryStatus::AuthenticationStopped { status, reason });
return;
}
InventoryDelivery::Rejected { status, reason } => {
let suffix = reason.map_or_else(String::new, |reason| format!("; reason={reason}"));
let _ = status_tx.send(InventoryStatus::Failed {
reason: format!("Connect rejected inventory with HTTP {status}{suffix}"),
});
return;
}
};
if sleep_or_cancel(&task_shutdown, delay).await {
break;
}
}
let _ = status_tx.send(InventoryStatus::Stopped);
});
Ok(Some(InventoryRuntime {
shutdown,
status: status_rx,
task: Some(task),
}))
}
@@ -289,12 +131,6 @@ fn failed(status: &watch::Sender<HeartbeatStatus>, error: HeartbeatError) {
});
}
fn failed_inventory(status: &watch::Sender<InventoryStatus>, error: InventoryError) {
let _ = status.send(InventoryStatus::Failed {
reason: error.to_string(),
});
}
fn jitter(maximum: Duration) -> Duration {
if maximum.is_zero() {
Duration::ZERO
-265
View File
@@ -1,265 +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 std::time::Duration;
use chrono::{DateTime, SecondsFormat, Utc};
use reqwest::{Client, StatusCode, Url, header};
use rustls::RootCertStore;
use rustls::pki_types::{CertificateDer, pem::PemObject as _};
use serde::{Deserialize, Serialize};
use zeroize::Zeroizing;
use super::config::HeartbeatConfig;
use super::credential_store::{CredentialStoreError, DeviceCredential};
use super::identity::IdentityError;
use super::identity_store::StoreError;
use super::registration::{CredentialValidationError, validate_stored_credential};
const MAX_RESPONSE_BYTES: usize = 64 * 1024;
pub(crate) enum TelemetryDelivery {
Accepted { cluster_name: String, body: Vec<u8> },
Retry { retry_after: Option<Duration> },
AuthenticationStopped { status: u16, reason: Option<String> },
Rejected { status: u16, reason: Option<String> },
}
pub(crate) struct TelemetryTransport {
endpoint: Url,
root_store: RootCertStore,
roots: Vec<CertificateDer<'static>>,
config: HeartbeatConfig,
}
impl TelemetryTransport {
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, TelemetryError> {
let mut endpoint = Url::parse(&config.endpoint).map_err(|_| TelemetryError::Endpoint)?;
if endpoint.scheme() != "https"
|| endpoint.cannot_be_a_base()
|| !endpoint.username().is_empty()
|| endpoint.password().is_some()
|| endpoint.query().is_some()
|| endpoint.fragment().is_some()
{
return Err(TelemetryError::Endpoint);
}
if !endpoint.path().ends_with('/') {
endpoint.set_path(&format!("{}/", endpoint.path()));
}
let roots = CertificateDer::pem_slice_iter(&config.root_ca_pem)
.collect::<Result<Vec<_>, _>>()
.map_err(|_| TelemetryError::RootCertificate)?;
if roots.is_empty() {
return Err(TelemetryError::RootCertificate);
}
let mut root_store = RootCertStore::empty();
let (accepted, rejected) = root_store.add_parsable_certificates(roots.clone());
if accepted != roots.len() || rejected != 0 {
return Err(TelemetryError::RootCertificate);
}
let schedule = config.schedule;
if schedule.timeout.is_zero()
|| schedule.timeout > Duration::from_secs(5)
|| schedule.initial_backoff.is_zero()
|| schedule.max_backoff < schedule.initial_backoff
|| schedule.max_backoff > Duration::from_secs(5 * 60)
{
return Err(TelemetryError::Schedule);
}
Ok(Self {
endpoint,
root_store,
roots,
config,
})
}
pub(crate) async fn post<T: Serialize>(&self, collection: &str, value: &T) -> Result<TelemetryDelivery, TelemetryError> {
let (cluster_name, cluster_uid, client) = self.authenticated_client().await?;
let url = self.endpoint.join(&format!("clusters/{cluster_uid}/{collection}"))?;
let response = match client.post(url).json(value).send().await {
Ok(response) => response,
Err(error) if error.is_timeout() || error.is_connect() || error.is_request() => {
return Ok(TelemetryDelivery::Retry { retry_after: None });
}
Err(error) => return Err(error.into()),
};
let status = response.status();
if status == StatusCode::TOO_MANY_REQUESTS {
return Ok(TelemetryDelivery::Retry {
retry_after: retry_after(response.headers(), Utc::now(), self.config.schedule.max_backoff),
});
}
if status == StatusCode::REQUEST_TIMEOUT || status.is_server_error() {
return Ok(TelemetryDelivery::Retry { retry_after: None });
}
if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) {
return Ok(TelemetryDelivery::AuthenticationStopped {
status: status.as_u16(),
reason: response_reason(response).await,
});
}
if status != StatusCode::OK {
return Ok(TelemetryDelivery::Rejected {
status: status.as_u16(),
reason: response_reason(response).await,
});
}
Ok(TelemetryDelivery::Accepted {
cluster_name,
body: bounded_body(response).await?,
})
}
async fn authenticated_client(&self) -> Result<(String, String, Client), TelemetryError> {
let _lock = self.config.credential_store.lock().await?;
let credential = self.config.credential_store.load()?.ok_or(TelemetryError::NotRegistered)?;
let identity = self.config.identity_store.load()?.ok_or(TelemetryError::IdentityMissing)?;
validate_stored_credential(&credential, &identity, &self.root_store, &self.roots)?;
let now = Utc::now().timestamp();
if now < credential.not_before_unix || now >= credential.not_after_unix {
return Err(TelemetryError::CredentialExpired);
}
let (organization_uid, cluster_uid) = credential_parent(&credential)?;
let cluster_name = format!("organizations/{organization_uid}/clusters/{cluster_uid}");
let client = self.client(&credential, &identity.to_pkcs8_pem()?)?;
Ok((cluster_name, cluster_uid.to_owned(), client))
}
fn client(&self, credential: &DeviceCredential, key: &Zeroizing<String>) -> Result<Client, TelemetryError> {
let mut pem = Zeroizing::new(Vec::with_capacity(credential.certificate_chain.len() + key.len() + 1));
pem.extend_from_slice(credential.certificate_chain.as_bytes());
pem.push(b'\n');
pem.extend_from_slice(key.as_bytes());
let identity = reqwest::Identity::from_pem(&pem).map_err(|_| TelemetryError::IdentityCertificate)?;
let roots = self
.roots
.iter()
.map(|root| reqwest::Certificate::from_der(root.as_ref()))
.collect::<Result<Vec<_>, _>>()?;
Client::builder()
.https_only(true)
.redirect(reqwest::redirect::Policy::none())
.timeout(self.config.schedule.timeout)
.tls_certs_only(roots)
.identity(identity)
.build()
.map_err(Into::into)
}
}
fn credential_parent(credential: &DeviceCredential) -> Result<(&str, &str), TelemetryError> {
let mut parts = credential.name.split('/');
let valid = parts.next() == Some("organizations");
let organization_uid = parts.next();
let valid = valid && parts.next() == Some("clusters");
let cluster_uid = parts.next();
let valid = valid && parts.next() == Some("clusterDevices");
let device_uid = parts.next();
if !valid
|| organization_uid.is_none_or(str::is_empty)
|| cluster_uid.is_none_or(str::is_empty)
|| device_uid != Some(credential.uid.as_str())
|| parts.next().is_some()
{
return Err(TelemetryError::CredentialName);
}
Ok((
organization_uid.ok_or(TelemetryError::CredentialName)?,
cluster_uid.ok_or(TelemetryError::CredentialName)?,
))
}
fn retry_after(headers: &header::HeaderMap, now: DateTime<Utc>, maximum: Duration) -> Option<Duration> {
let value = headers.get(header::RETRY_AFTER)?.to_str().ok()?;
let delay = value.parse::<u64>().ok().map(Duration::from_secs).or_else(|| {
DateTime::parse_from_rfc2822(value)
.ok()
.and_then(|at| (at.with_timezone(&Utc) - now).to_std().ok())
})?;
Some(delay.min(maximum))
}
pub(crate) fn is_exact_utc_seconds(value: &str) -> bool {
DateTime::parse_from_rfc3339(value).is_ok_and(|time| {
time.offset().local_minus_utc() == 0
&& value.ends_with('Z')
&& time.with_timezone(&Utc).to_rfc3339_opts(SecondsFormat::Secs, true) == value
})
}
async fn response_reason(response: reqwest::Response) -> Option<String> {
#[derive(Deserialize)]
struct Envelope {
#[serde(default)]
details: Vec<Detail>,
}
#[derive(Deserialize)]
struct Detail {
#[serde(default)]
reason: String,
}
serde_json::from_slice::<Envelope>(&bounded_body(response).await.ok()?)
.ok()?
.details
.into_iter()
.find_map(|detail| (!detail.reason.is_empty()).then_some(detail.reason))
}
async fn bounded_body(mut response: reqwest::Response) -> Result<Vec<u8>, TelemetryError> {
let mut body = Vec::new();
while let Some(chunk) = response.chunk().await? {
if body.len().saturating_add(chunk.len()) > MAX_RESPONSE_BYTES {
return Err(TelemetryError::ResponseTooLarge);
}
body.extend_from_slice(&chunk);
}
Ok(body)
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum TelemetryError {
#[error("Connect telemetry endpoint must be an HTTPS base URL without credentials, query, or fragment")]
Endpoint,
#[error("Connect telemetry root CA configuration is invalid")]
RootCertificate,
#[error("Connect telemetry retry schedule is invalid")]
Schedule,
#[error("RustFS is not registered with Connect")]
NotRegistered,
#[error("the Connect device private key is missing")]
IdentityMissing,
#[error("the stored Connect certificate and device private key cannot form a TLS identity")]
IdentityCertificate,
#[error("the stored Connect credential name is invalid")]
CredentialName,
#[error("the stored Connect device certificate is not currently valid")]
CredentialExpired,
#[error("Connect telemetry response exceeded 64 KiB")]
ResponseTooLarge,
#[error(transparent)]
Url(#[from] url::ParseError),
#[error(transparent)]
Transport(#[from] reqwest::Error),
#[error(transparent)]
Identity(#[from] IdentityError),
#[error(transparent)]
IdentityStore(#[from] StoreError),
#[error(transparent)]
CredentialStore(#[from] CredentialStoreError),
#[error(transparent)]
CredentialValidation(#[from] CredentialValidationError),
}
+10 -8
View File
@@ -26,22 +26,22 @@ struct MiMallocAllocator;
unsafe impl GlobalAlloc for MiMallocAllocator {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { mimalloc::MiMalloc.alloc(layout) }
unsafe { rustfs_mimalloc::MiMalloc.alloc(layout) }
}
unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { mimalloc::MiMalloc.alloc_zeroed(layout) }
unsafe { rustfs_mimalloc::MiMalloc.alloc_zeroed(layout) }
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { mimalloc::MiMalloc.dealloc(ptr, layout) }
unsafe { rustfs_mimalloc::MiMalloc.dealloc(ptr, layout) }
}
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
unsafe { rustfs_mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
}
}
@@ -51,7 +51,7 @@ static GLOBAL: hotpath::CountingAllocator<MiMallocAllocator> = hotpath::Counting
#[cfg(not(all(feature = "hotpath", feature = "hotpath-alloc")))]
#[global_allocator]
static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc;
static GLOBAL: rustfs_mimalloc::MiMalloc = rustfs_mimalloc::MiMalloc;
fn main() {
let _hotpath_guard = hotpath::HotpathGuardBuilder::new("main").build();
@@ -71,8 +71,9 @@ mod tests {
allocation.extend_from_slice(&[7_u8; 64]);
assert_eq!(allocation.len(), 64);
let heap = rustfs_mimalloc::heap::Heap::main();
// SAFETY: the live Vec pointer is valid to inspect for heap ownership.
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(allocation.as_ptr().cast()) });
assert!(unsafe { heap.contains(allocation.as_ptr()) });
}
#[test]
@@ -85,12 +86,13 @@ mod tests {
let layout = Layout::from_size_align(32, 8).expect("valid test allocation layout");
let grown_layout = Layout::from_size_align(64, 8).expect("valid grown test allocation layout");
let allocator = super::MiMallocAllocator;
let heap = rustfs_mimalloc::heap::Heap::main();
// SAFETY: The pointer is checked for null before use and later released
// through the same allocator with the corresponding layout.
let ptr = unsafe { allocator.alloc_zeroed(layout) };
assert!(!ptr.is_null());
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(ptr.cast()) });
assert!(unsafe { heap.contains(ptr) });
assert!(unsafe { std::slice::from_raw_parts(ptr, 32).iter().all(|byte| *byte == 0) });
// SAFETY: `ptr` was allocated by `allocator` with `layout`; on failure
@@ -102,7 +104,7 @@ mod tests {
panic!("mimalloc realloc failed in allocator smoke test");
}
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(grown_ptr.cast()) });
assert!(unsafe { heap.contains(grown_ptr) });
// SAFETY: `grown_ptr` was reallocated by `allocator` and is released
// with the matching grown layout.
unsafe { allocator.dealloc(grown_ptr, grown_layout) };
+19 -36
View File
@@ -17,10 +17,7 @@ use rustfs_io_metrics::{
record_cpu_usage, record_memory_usage, record_process_memory_split,
};
use serde::Serialize;
#[cfg(any(test, not(target_os = "windows")))]
use serde_json::Value;
#[cfg(not(target_os = "windows"))]
use std::ffi::CStr;
use std::path::Path;
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
@@ -231,7 +228,18 @@ fn read_cgroup_memory_snapshot() -> Option<CgroupMemorySnapshot> {
read_cgroup_v2().or_else(read_cgroup_v1)
}
#[cfg(any(test, not(target_os = "windows")))]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
let json = rustfs_mimalloc::MiMalloc::stats_json();
if json.is_empty() {
return None;
}
let observation = parse_mimalloc_stats_json(&json)?;
Some(AllocatorMemorySnapshot {
backend: crate::allocator_reclaim::allocator_backend(),
observation,
})
}
fn numeric_json_value(value: &Value) -> Option<u64> {
match value {
Value::Number(number) => number
@@ -242,7 +250,6 @@ fn numeric_json_value(value: &Value) -> Option<u64> {
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
match value {
Value::Object(fields) => fields
@@ -254,7 +261,6 @@ fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64> {
match value {
Value::Object(fields) => {
@@ -271,12 +277,10 @@ fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64>
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_current(value: &Value, metric: &str) -> Option<u64> {
mimalloc_stat_field(value, metric, "current")
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64> {
metrics
.iter()
@@ -285,7 +289,6 @@ fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64
.filter(|value| *value > 0)
}
#[cfg(any(test, not(target_os = "windows")))]
fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservation> {
let value = serde_json::from_str::<Value>(stats_json).ok()?;
let malloc_metrics = ["malloc_normal", "malloc_huge"];
@@ -312,33 +315,6 @@ fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservat
}
}
#[cfg(not(target_os = "windows"))]
#[allow(unsafe_code)]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
// SAFETY: `mi_stats_get_json` returns a null-terminated JSON buffer owned by
// mimalloc when called with a null input buffer. The mimalloc API requires
// freeing that buffer with `mi_free`; parsing finishes before the buffer is freed.
let observation = unsafe {
let stats_ptr = libmimalloc_sys::mi_stats_get_json(0, std::ptr::null_mut());
if stats_ptr.is_null() {
return None;
}
let observation = CStr::from_ptr(stats_ptr).to_str().ok().and_then(parse_mimalloc_stats_json);
libmimalloc_sys::mi_free(stats_ptr.cast());
observation?
};
Some(AllocatorMemorySnapshot {
backend: crate::allocator_reclaim::allocator_backend(),
observation,
})
}
#[cfg(target_os = "windows")]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
None
}
fn configured_memory_observability_interval_secs() -> u64 {
rustfs_utils::get_env_u64(ENV_MEMORY_OBSERVABILITY_INTERVAL_SECS, DEFAULT_MEMORY_OBSERVABILITY_INTERVAL_SECS).max(1)
}
@@ -566,6 +542,13 @@ mod tests {
assert_eq!(parse_mimalloc_stats_json(r#"{ "allocator": "unknown" }"#), None);
}
#[test]
fn read_allocator_memory_snapshot_uses_mimalloc_stats_json() {
let snapshot = super::read_allocator_memory_snapshot();
#[cfg(not(target_os = "windows"))]
assert!(snapshot.is_some(), "allocator snapshot should be available on non-Windows");
}
#[test]
fn memory_observability_snapshot_reports_disabled_when_metrics_are_disabled() {
let snapshot = build_memory_observability_status_snapshot(false, 15, false);
+3 -77
View File
@@ -13,13 +13,10 @@
// limitations under the License.
use crate::site_replication_reconcile::spawn_site_replication_reconcile_task;
use crate::storage_api::startup::services::{ECStore, EndpointServerPools, ServerContextSlot, StorageAdminApi};
use crate::storage_api::startup::services::{ECStore, EndpointServerPools, ServerContextSlot};
use crate::{
config::Config,
connect::{
CoarseNodeSummary, HeartbeatConfig, HeartbeatRuntime, InventoryError, InventoryFlag, InventoryRuntime, InventorySchedule,
InventorySnapshot, spawn_heartbeat_runtime, spawn_inventory_runtime,
},
connect::{CoarseNodeSummary, HeartbeatConfig, HeartbeatRuntime, spawn_heartbeat_runtime},
init::{init_buffer_profile_system, init_kms_system},
server::ServiceStateManager,
startup_audit::init_audit_runtime,
@@ -80,9 +77,6 @@ pub(crate) async fn init_startup_runtime_services(
let optional_runtimes = init_optional_runtime_services().await?;
let heartbeat_config = HeartbeatConfig::from_env().map_err(std::io::Error::other)?;
let heartbeat_nodes = heartbeat_config.as_ref().map(|_| endpoint_pools.get_nodes().len());
let inventory_drives = heartbeat_config
.as_ref()
.map(|_| endpoint_pools.as_ref().iter().map(|pool| pool.endpoints.as_ref().len()).sum());
init_buffer_profile_system(config);
init_deadlock_detector_runtime();
@@ -102,9 +96,7 @@ pub(crate) async fn init_startup_runtime_services(
init_notification_runtime(endpoint_pools, buckets).await?;
let enable_scanner = init_background_service_runtime(store.clone()).await?;
init_observability_runtime(store.clone(), ctx.clone()).await;
let heartbeat = start_heartbeat_runtime(heartbeat_config.clone(), heartbeat_nodes, &ctx)?;
let inventory = start_inventory_runtime(heartbeat_config, heartbeat_nodes, inventory_drives, store, &ctx)?;
let heartbeat = heartbeat.map(|heartbeat| heartbeat.with_inventory(inventory));
let heartbeat = start_heartbeat_runtime(heartbeat_config, heartbeat_nodes, &ctx)?;
Ok(StartupServiceRuntime {
optional_runtimes,
@@ -128,69 +120,3 @@ fn start_heartbeat_runtime(
.ok_or_else(|| std::io::Error::other("Connect heartbeat node count is outside protocol bounds"))?;
spawn_heartbeat_runtime(Some(config), shutdown, move || summary).map_err(std::io::Error::other)
}
fn start_inventory_runtime(
config: Option<HeartbeatConfig>,
node_count: Option<usize>,
expected_drive_count: Option<usize>,
store: Arc<ECStore>,
shutdown: &CancellationToken,
) -> Result<Option<InventoryRuntime>> {
let Some(config) = config else {
return Ok(None);
};
let node_count = node_count.unwrap_or_default();
let expected_drive_count = expected_drive_count.unwrap_or_default();
spawn_inventory_runtime(Some(config), InventorySchedule::default(), shutdown, move || {
let store = store.clone();
async move {
let info = StorageAdminApi::storage_info(store.as_ref()).await;
inventory_snapshot(node_count, expected_drive_count, info)
}
})
.map_err(std::io::Error::other)
}
fn inventory_snapshot(
node_count: usize,
expected_drive_count: usize,
info: rustfs_madmin::StorageInfo,
) -> std::result::Result<InventorySnapshot, InventoryError> {
if info.disks.len() != expected_drive_count {
return Err(InventoryError::SnapshotIncomplete {
expected: expected_drive_count,
observed: info.disks.len(),
});
}
let total = crate::app::storage_api::capacity::get_total_usable_capacity(&info.disks, &info) as u64;
let free = crate::app::storage_api::capacity::get_total_usable_capacity_free(&info.disks, &info) as u64;
let mut flags = Vec::with_capacity(3);
if info.disks.iter().any(|disk| disk.state == rustfs_madmin::ITEM_OFFLINE) {
flags.extend([InventoryFlag::ClusterDegraded, InventoryFlag::DriveOffline]);
}
if info.disks.iter().any(|disk| disk.healing) {
flags.push(InventoryFlag::ClusterHealing);
}
InventorySnapshot::current(node_count, info.disks.len(), total, free, flags)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn inventory_rejects_a_partial_startup_storage_snapshot() {
let info = rustfs_madmin::StorageInfo {
disks: vec![rustfs_madmin::Disk::default()],
..Default::default()
};
assert!(matches!(
inventory_snapshot(2, 2, info),
Err(InventoryError::SnapshotIncomplete {
expected: 2,
observed: 1
})
));
}
}
-1
View File
@@ -279,7 +279,6 @@ pub(crate) mod startup {
}
pub(crate) mod services {
pub(crate) use super::super::storage_contracts::StorageAdminApi;
pub(crate) use crate::storage::storage_api::{ECStore, EndpointServerPools, ServerContextSlot};
}
-669
View File
@@ -1,669 +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 std::collections::VecDeque;
use std::fs;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use bytes::Bytes;
use http_body_util::{BodyExt as _, Full};
use hyper::service::service_fn;
use hyper::{Request, Response, StatusCode};
use hyper_util::rt::TokioIo;
use rcgen::{
BasicConstraints, CertificateParams, DistinguishedName, DnType, ExtendedKeyUsagePurpose, IsCa, Issuer, KeyPair,
KeyUsagePurpose, SanType,
};
use rustfs::connect::{
CredentialStore, DeviceCredential, HeartbeatConfig, HeartbeatSchedule, IdentityStore, InventoryFlag, InventoryOsVersion,
InventorySchedule, InventorySnapshot, InventoryStatus, OperatingSystemFamily, spawn_inventory_runtime,
};
use rustls::RootCertStore;
use rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer};
use rustls::server::WebPkiClientVerifier;
use serde_json::{Value, json};
use time::OffsetDateTime;
use tokio::net::TcpListener;
use tokio::sync::watch;
use tokio_rustls::TlsAcceptor;
use tokio_util::sync::CancellationToken;
const ORGANIZATION_UID: &str = "0198f4b0-1a00-7c10-8d21-2e3f4a5b6c70";
const CLUSTER_UID: &str = "0198f4b0-2b00-7d20-9e31-3f4a5b6c7d81";
const DEVICE_UID: &str = "0198f4b0-3c00-7e30-8f41-4a5b6c7d8e92";
const SNAPSHOT_UID: &str = "0198f4b0-4d00-7f40-9051-5b6c7d8e9fa3";
struct TestPki {
root_params: CertificateParams,
root_key: KeyPair,
root_der: CertificateDer<'static>,
root_pem: String,
server_der: CertificateDer<'static>,
server_key: PrivatePkcs8KeyDer<'static>,
}
impl TestPki {
fn new() -> Self {
let now = OffsetDateTime::now_utc();
let root_key = KeyPair::generate().expect("generate root key");
let mut root_params = CertificateParams::default();
root_params.is_ca = IsCa::Ca(BasicConstraints::Unconstrained);
root_params.not_before = now - time::Duration::days(30);
root_params.not_after = now + time::Duration::days(30);
root_params.key_usages = vec![KeyUsagePurpose::KeyCertSign, KeyUsagePurpose::DigitalSignature];
let root = root_params.self_signed(&root_key).expect("sign root");
let server_key = KeyPair::generate().expect("generate server key");
let mut server_params = CertificateParams::default();
server_params.not_before = now - time::Duration::hours(1);
server_params.not_after = now + time::Duration::days(2);
server_params
.subject_alt_names
.push(SanType::DnsName("localhost".try_into().expect("valid DNS name")));
server_params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ServerAuth];
let server = server_params
.signed_by(&server_key, &Issuer::from_params(&root_params, &root_key))
.expect("sign server certificate");
Self {
root_params,
root_key,
root_der: root.der().clone(),
root_pem: root.pem(),
server_der: server.der().clone(),
server_key: PrivatePkcs8KeyDer::from(server_key.serialize_der()),
}
}
fn server_config(&self) -> rustls::ServerConfig {
let mut roots = RootCertStore::empty();
roots.add(self.root_der.clone()).expect("add client root");
let verifier = WebPkiClientVerifier::builder(Arc::new(roots))
.build()
.expect("client verifier");
rustls::ServerConfig::builder()
.with_client_cert_verifier(verifier)
.with_single_cert(vec![self.server_der.clone()], PrivateKeyDer::Pkcs8(self.server_key.clone_key()))
.expect("server TLS")
}
fn stores(&self, temp: &tempfile::TempDir) -> (IdentityStore, CredentialStore) {
let identity_store = IdentityStore::new(temp.path().join("identity"));
let identity = identity_store.load_or_create().expect("create identity");
let private_key = PrivatePkcs8KeyDer::from(identity.to_pkcs8_der().expect("serialize key").to_vec());
let device_key = KeyPair::from_pkcs8_der_and_sign_algo(&private_key, &rcgen::PKCS_ECDSA_P256_SHA256).expect("device key");
let now = OffsetDateTime::now_utc();
let mut params = CertificateParams::default();
params.not_before = now - time::Duration::hours(1);
params.not_after = now + time::Duration::hours(23);
params.serial_number = Some(vec![1; 16].into());
params.key_usages = vec![KeyUsagePurpose::DigitalSignature];
params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ClientAuth];
params.distinguished_name = DistinguishedName::new();
params.distinguished_name.push(DnType::CommonName, DEVICE_UID);
params.subject_alt_names.push(SanType::URI(
format!("urn:rustfs:connect:device:{DEVICE_UID}")
.try_into()
.expect("device URI"),
));
let certificate = params
.signed_by(&device_key, &Issuer::from_params(&self.root_params, &self.root_key))
.expect("device certificate");
let cluster = format!("organizations/{ORGANIZATION_UID}/clusters/{CLUSTER_UID}");
let credential = DeviceCredential {
name: format!("{cluster}/clusterDevices/{DEVICE_UID}"),
uid: DEVICE_UID.to_owned(),
protocol_version: "v1".to_owned(),
key_id: format!("x509-{}", "01".repeat(16)),
certificate_serial: "01".repeat(16),
certificate: certificate.pem(),
certificate_chain: certificate.pem(),
not_before_unix: (now - time::Duration::hours(1)).unix_timestamp(),
not_after_unix: (now + time::Duration::hours(23)).unix_timestamp(),
};
let directory = temp.path().join("credential");
fs::create_dir_all(&directory).expect("credential directory");
let path = directory.join("device.crt.json");
fs::write(&path, serde_json::to_vec(&credential).expect("credential JSON")).expect("write credential");
private_mode(&path);
(identity_store, CredentialStore::new(directory))
}
}
#[derive(Clone)]
struct Reply {
status: StatusCode,
body: Value,
retry_after: Option<&'static str>,
}
impl Reply {
fn ok(content_hash: &str) -> Self {
Self {
status: StatusCode::OK,
body: json!({
"name": format!("organizations/{ORGANIZATION_UID}/clusters/{CLUSTER_UID}/inventorySnapshots/{SNAPSHOT_UID}"),
"uid": SNAPSHOT_UID,
"contentHash": content_hash,
"receivedAt": "2026-08-22T01:02:03Z",
"futureField": true
}),
retry_after: None,
}
}
fn error(status: StatusCode, reason: &str) -> Self {
Self {
status,
body: json!({"details": [{"reason": reason}]}),
retry_after: None,
}
}
}
struct TestServer {
endpoint: String,
seen: Arc<Mutex<Vec<Value>>>,
task: tokio::task::JoinHandle<()>,
}
impl Drop for TestServer {
fn drop(&mut self) {
self.task.abort();
}
}
async fn server(pki: &TestPki, replies: Vec<Reply>) -> TestServer {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind server");
let address = listener.local_addr().expect("server address");
let acceptor = TlsAcceptor::from(Arc::new(pki.server_config()));
let replies = Arc::new(Mutex::new(VecDeque::from(replies)));
let seen = Arc::new(Mutex::new(Vec::new()));
let captured = seen.clone();
let task = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let acceptor = acceptor.clone();
let replies = replies.clone();
let seen = captured.clone();
tokio::spawn(async move {
let Ok(stream) = acceptor.accept(stream).await else { return };
let service = service_fn(move |request: Request<hyper::body::Incoming>| {
let replies = replies.clone();
let seen = seen.clone();
async move {
assert_eq!(request.uri().path(), format!("/agent/clusters/{CLUSTER_UID}/inventorySnapshots"));
let body = request.into_body().collect().await.expect("request body").to_bytes();
seen.lock()
.expect("seen lock")
.push(serde_json::from_slice(&body).expect("request JSON"));
let reply = replies
.lock()
.expect("reply lock")
.pop_front()
.unwrap_or_else(|| Reply::error(StatusCode::SERVICE_UNAVAILABLE, "UNAVAILABLE"));
let mut builder = Response::builder()
.status(reply.status)
.header("content-type", "application/json");
if let Some(value) = reply.retry_after {
builder = builder.header("retry-after", value);
}
Ok::<_, hyper::Error>(
builder
.body(Full::new(Bytes::from(serde_json::to_vec(&reply.body).expect("reply JSON"))))
.expect("reply"),
)
}
});
let _ = hyper::server::conn::http1::Builder::new()
.serve_connection(TokioIo::new(stream), service)
.await;
});
}
});
TestServer {
endpoint: format!("https://localhost:{}/agent/", address.port()),
seen,
task,
}
}
fn config(temp: &tempfile::TempDir, pki: &TestPki, server: &TestServer) -> HeartbeatConfig {
let (identity_store, credential_store) = pki.stores(temp);
HeartbeatConfig {
endpoint: server.endpoint.clone(),
root_ca_pem: pki.root_pem.as_bytes().to_vec(),
identity_store,
credential_store,
state_path: temp.path().join("private-config-secret/heartbeat/state.json"),
schedule: HeartbeatSchedule {
cadence: Duration::from_secs(30),
jitter: Duration::ZERO,
timeout: Duration::from_millis(200),
initial_backoff: Duration::from_millis(20),
max_backoff: Duration::from_millis(80),
},
}
}
fn schedule() -> InventorySchedule {
InventorySchedule {
cadence: Duration::from_secs(60),
jitter: Duration::ZERO,
}
}
fn snapshot() -> InventorySnapshot {
InventorySnapshot::new(
"1.4.2",
Some(InventoryOsVersion::new(OperatingSystemFamily::Linux, 6, 8).expect("valid operating-system version")),
8,
96,
1_099_511_627_776,
412_316_860_416,
[InventoryFlag::ClusterDegraded, InventoryFlag::DriveOffline],
)
.expect("valid inventory")
}
fn collect_strings(value: &Value, strings: &mut Vec<String>) {
match value {
Value::String(value) => strings.push(value.clone()),
Value::Array(values) => values.iter().for_each(|value| collect_strings(value, strings)),
Value::Object(values) => values.values().for_each(|value| collect_strings(value, strings)),
_ => {}
}
}
async fn wait_for(
status: &mut watch::Receiver<InventoryStatus>,
predicate: impl Fn(&InventoryStatus) -> bool,
) -> InventoryStatus {
tokio::time::timeout(Duration::from_secs(3), async {
loop {
let current = status.borrow_and_update().clone();
if predicate(&current) {
return current;
}
status.changed().await.expect("status channel");
}
})
.await
.expect("inventory status timeout")
}
#[test]
fn connect_inventory_frozen_vector_has_the_exact_canonical_hash_and_no_open_ended_fields() {
let fixtures: Value = serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/inventory/valid-vectors.json"))
.expect("valid fixture JSON");
let expected = &fixtures["vectors"][0]["expected"];
let snapshot = snapshot();
assert_eq!(snapshot.content_hash().expect("content hash"), expected["contentHash"]);
assert_eq!(InventorySchedule::default().cadence, Duration::from_secs(6 * 60 * 60));
assert_eq!(InventorySchedule::default().jitter, Duration::from_secs(30 * 60));
let encoded = serde_json::to_value(snapshot).expect("snapshot JSON");
assert_eq!(
encoded,
json!({
"rustfsVersion": "1.4.2",
"osVersion": {"family": "linux", "major": 6, "minor": 8},
"nodeCount": 8,
"driveCount": 96,
"capacityTotalBytes": 1099511627776_u64,
"capacityUsedBytes": 412316860416_u64,
"coarseFlags": ["cluster.degraded", "drive.offline"]
})
);
let fixtures: Value =
serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/inventory/secret-like-vectors.json"))
.expect("valid secret-like fixture JSON");
let known_fields = [
"protocolVersion",
"rustfsVersion",
"osVersion",
"nodeCount",
"driveCount",
"capacityTotalBytes",
"capacityUsedBytes",
"coarseFlags",
];
let known_flags = ["cluster.degraded", "drive.offline"];
let mut excluded = Vec::new();
for vector in fixtures["vectors"].as_array().expect("fixture vectors") {
let input = vector["input"].as_object().expect("fixture input");
for (name, value) in input {
if !known_fields.contains(&name.as_str()) {
collect_strings(value, &mut excluded);
}
}
for (name, value) in input["osVersion"].as_object().expect("fixture OS version") {
if !["family", "major", "minor"].contains(&name.as_str()) {
collect_strings(value, &mut excluded);
}
}
for flag in input["coarseFlags"].as_array().expect("fixture coarse flags") {
let flag = flag.as_str().expect("fixture coarse flag");
if !known_flags.contains(&flag) {
excluded.push(flag.to_owned());
}
}
}
let encoded = serde_json::to_string(&encoded).expect("encoded snapshot");
for value in excluded {
assert!(!encoded.contains(&value), "snapshot exposed fixture value {value}");
}
}
#[test]
fn connect_inventory_bounds_fail_instead_of_truncating_or_inventing_values() {
assert!(matches!(
InventorySnapshot::current(0, 0, 0, 0, []),
Err(rustfs::connect::InventoryError::NodeCount)
));
assert!(matches!(
InventorySnapshot::current(1, 1_048_577, 0, 0, []),
Err(rustfs::connect::InventoryError::DriveCount)
));
assert!(matches!(
InventorySnapshot::current(1, 0, 9_007_199_254_740_992, 0, []),
Err(rustfs::connect::InventoryError::Capacity)
));
assert!(matches!(
InventorySnapshot::current(1, 0, 10, 11, []),
Err(rustfs::connect::InventoryError::Capacity)
));
assert!(matches!(
InventorySnapshot::new("1.0.0-private.1", None, 1, 0, 0, 0, []),
Err(rustfs::connect::InventoryError::RustfsVersion)
));
}
#[tokio::test]
async fn connect_inventory_restart_replays_the_pending_request_and_then_skips_unchanged_inventory() {
let pki = TestPki::new();
let content_hash = snapshot().content_hash().expect("content hash");
let first_server = server(&pki, vec![Reply::error(StatusCode::SERVICE_UNAVAILABLE, "UNAVAILABLE")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let samples = Arc::new(AtomicUsize::new(0));
let sampled = samples.clone();
let runtime = spawn_inventory_runtime(Some(config(&temp, &pki, &first_server)), schedule(), &shutdown, move || {
sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(Ok(snapshot()))
})
.expect("start inventory")
.expect("configured inventory");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::BackingOff { .. })).await,
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(20)
));
assert_eq!(samples.load(Ordering::Relaxed), 1);
let original = first_server.seen.lock().expect("seen lock")[0].clone();
runtime.shutdown().await;
let mut limited = Reply::error(StatusCode::TOO_MANY_REQUESTS, "RATE_LIMITED");
limited.retry_after = Some("0");
let restart_server = server(&pki, vec![limited, Reply::ok(&content_hash)]).await;
let restart_config = config(&temp, &pki, &restart_server);
let restart_samples = Arc::new(AtomicUsize::new(0));
let sampled = restart_samples.clone();
let restart = spawn_inventory_runtime(Some(restart_config.clone()), schedule(), &shutdown, move || {
sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(Ok(snapshot()))
})
.expect("restart inventory")
.expect("configured inventory");
let mut status = restart.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::Online { .. })).await,
InventoryStatus::Online { content_hash: accepted, received_at }
if accepted == content_hash && received_at == "2026-08-22T01:02:03Z"
));
assert_eq!(restart_samples.load(Ordering::Relaxed), 0);
let delivered = restart_server.seen.lock().expect("seen lock").clone();
assert_eq!(delivered, vec![original.clone(), original.clone()]);
assert_eq!(original["sequence"], 0);
let encoded = serde_json::to_string(&original).expect("request JSON");
for forbidden in [
"private-config-secret",
"BEGIN CERTIFICATE",
"AKIAIOSFODNN7EXAMPLE",
"bucket",
"object",
"path",
] {
assert!(!encoded.contains(forbidden), "request exposed {forbidden}");
}
assert_eq!(original.as_object().expect("request object").len(), 10);
restart.shutdown().await;
let unchanged_samples = Arc::new(AtomicUsize::new(0));
let sampled = unchanged_samples.clone();
let unchanged = spawn_inventory_runtime(Some(restart_config), schedule(), &shutdown, move || {
sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(Ok(snapshot()))
})
.expect("restart inventory")
.expect("configured inventory");
let mut unchanged_status = unchanged.status();
assert!(matches!(
wait_for(&mut unchanged_status, |status| matches!(status, InventoryStatus::Unchanged { .. })).await,
InventoryStatus::Unchanged { content_hash: unchanged } if unchanged == content_hash
));
assert_eq!(unchanged_samples.load(Ordering::Relaxed), 1);
assert_eq!(restart_server.seen.lock().expect("seen lock").len(), 2);
unchanged.shutdown().await;
}
#[tokio::test]
async fn connect_inventory_disconnect_retries_without_resampling() {
let pki = TestPki::new();
let unavailable = server(&pki, Vec::new()).await;
let temp = tempfile::tempdir().expect("tempdir");
let config = config(&temp, &pki, &unavailable);
drop(unavailable);
let shutdown = CancellationToken::new();
let samples = Arc::new(AtomicUsize::new(0));
let sampled = samples.clone();
let runtime = spawn_inventory_runtime(Some(config), schedule(), &shutdown, move || {
sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(Ok(snapshot()))
})
.expect("start inventory")
.expect("configured inventory");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| {
matches!(status, InventoryStatus::BackingOff { delay } if *delay == Duration::from_millis(40))
})
.await,
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(40)
));
assert_eq!(samples.load(Ordering::Relaxed), 1);
runtime.shutdown().await;
}
#[tokio::test]
async fn connect_inventory_retries_an_incomplete_sample_before_delivery() {
let pki = TestPki::new();
let content_hash = snapshot().content_hash().expect("content hash");
let server = server(&pki, vec![Reply::ok(&content_hash)]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let samples = Arc::new(AtomicUsize::new(0));
let sampled = samples.clone();
let runtime = spawn_inventory_runtime(Some(config(&temp, &pki, &server)), schedule(), &shutdown, move || {
let attempt = sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(if attempt == 0 {
Err(rustfs::connect::InventoryError::SnapshotIncomplete {
expected: 96,
observed: 12,
})
} else {
Ok(snapshot())
})
})
.expect("start inventory")
.expect("configured inventory");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::Online { .. })).await,
InventoryStatus::Online { content_hash: accepted, .. } if accepted == content_hash
));
assert_eq!(samples.load(Ordering::Relaxed), 2);
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
runtime.shutdown().await;
}
#[tokio::test]
async fn connect_inventory_unchanged_sample_resets_incomplete_backoff() {
let pki = TestPki::new();
let content_hash = snapshot().content_hash().expect("content hash");
let server = server(&pki, vec![Reply::ok(&content_hash)]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
let seed = spawn_inventory_runtime(Some(config.clone()), schedule(), &shutdown, || std::future::ready(Ok(snapshot())))
.expect("start inventory")
.expect("configured inventory");
let mut seed_status = seed.status();
assert!(matches!(
wait_for(&mut seed_status, |status| matches!(status, InventoryStatus::Online { .. })).await,
InventoryStatus::Online { content_hash: accepted, .. } if accepted == content_hash
));
seed.shutdown().await;
let samples = Arc::new(AtomicUsize::new(0));
let sampled = samples.clone();
let runtime = spawn_inventory_runtime(
Some(config),
InventorySchedule {
cadence: Duration::from_millis(100),
jitter: Duration::ZERO,
},
&shutdown,
move || {
let attempt = sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(if matches!(attempt, 0 | 1 | 3) {
Err(rustfs::connect::InventoryError::SnapshotIncomplete {
expected: 96,
observed: 12,
})
} else {
Ok(snapshot())
})
},
)
.expect("restart inventory")
.expect("configured inventory");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| {
matches!(status, InventoryStatus::BackingOff { delay } if *delay == Duration::from_millis(20))
})
.await,
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(20)
));
assert!(matches!(
wait_for(&mut status, |status| {
matches!(status, InventoryStatus::BackingOff { delay } if *delay == Duration::from_millis(40))
})
.await,
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(40)
));
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::Unchanged { .. })).await,
InventoryStatus::Unchanged { content_hash: unchanged } if unchanged == content_hash
));
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::BackingOff { .. })).await,
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(20)
));
assert_eq!(samples.load(Ordering::Relaxed), 4);
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
runtime.shutdown().await;
}
#[tokio::test]
async fn connect_inventory_revoked_device_stops_without_retrying() {
let pki = TestPki::new();
let server = server(&pki, vec![Reply::error(StatusCode::UNAUTHORIZED, "DEVICE_REVOKED")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let runtime = spawn_inventory_runtime(Some(config(&temp, &pki, &server)), schedule(), &shutdown, || {
std::future::ready(Ok(snapshot()))
})
.expect("start inventory")
.expect("configured inventory");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::AuthenticationStopped { .. })).await,
InventoryStatus::AuthenticationStopped { status: 401, reason: Some(reason) } if reason == "DEVICE_REVOKED"
));
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
runtime.shutdown().await;
}
#[tokio::test]
async fn connect_inventory_sequence_overflow_fails_before_sampling_or_network_delivery() {
let pki = TestPki::new();
let server = server(&pki, Vec::new()).await;
let temp = tempfile::tempdir().expect("tempdir");
let config = config(&temp, &pki, &server);
let state = temp.path().join("private-config-secret/inventory/state.json");
fs::create_dir_all(state.parent().expect("state directory")).expect("create state directory");
fs::write(
&state,
br#"{"nextSequence":9007199254740992,"pending":null,"lastAcceptedContentHash":null}"#,
)
.expect("write state");
private_mode(&state);
let samples = Arc::new(AtomicUsize::new(0));
let sampled = samples.clone();
let shutdown = CancellationToken::new();
let runtime = spawn_inventory_runtime(Some(config), schedule(), &shutdown, move || {
sampled.fetch_add(1, Ordering::Relaxed);
std::future::ready(Ok(snapshot()))
})
.expect("start inventory")
.expect("configured inventory");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, InventoryStatus::Failed { .. })).await,
InventoryStatus::Failed { reason } if reason.contains("sequence is exhausted")
));
assert_eq!(samples.load(Ordering::Relaxed), 0);
assert!(server.seen.lock().expect("seen lock").is_empty());
runtime.shutdown().await;
}
#[cfg(unix)]
fn private_mode(path: &std::path::Path) {
use std::os::unix::fs::PermissionsExt as _;
fs::set_permissions(path, fs::Permissions::from_mode(0o600)).expect("private permissions");
}
#[cfg(not(unix))]
fn private_mode(_path: &std::path::Path) {}
+1 -1
View File
@@ -241,7 +241,7 @@ env \
RUSTFS_TEST_VAULT_FAILOVER_MARKER="$MARKER" \
RUSTFS_TEST_VAULT_OLD_LEADER="$OLD_LEADER" \
cargo test -p rustfs-kms --test vault_ha_failover_live \
vault_raft_leader_failure_preserves_kv2_and_transit_decrypts -- \
vault_raft_leader_failure_recovers_kv2_and_transit_decrypts -- \
--ignored --nocapture --test-threads=1 &
TEST_PID=$!