fix: prevent metadata lock reentry and stabilize CI fixtures (#8257)

* fix(test): establish a writable previous-release upgrade baseline

* fix(ci): reserve capacity for durable admin fixtures

* fix(ci): group durable IAM state fixtures by resource needs

* ci: run E2E doctests with the E2E dependency graph

* fix(ci): separate fixture startup from transport deadlines

* fix(ci): bound pagination after seeding and revisit restored copies

* fix(ci): make filesystem fixture timing deterministic

* fix(ecstore): avoid metadata lock reentry during internal mutations

* test: align recovery fixtures with durable ownership contracts
This commit is contained in:
Chris
2026-09-30 20:57:07 +08:00
committed by GitHub
parent 8655c38f1d
commit 2806a80c91
19 changed files with 764 additions and 228 deletions
+1 -1
View File
@@ -1 +1 @@
sha256=7e233cd838efd2688a9cf11ba23ad60297172d124067e91f59502d10c9633464
sha256=476cb92aa9e385c9bb509ee5df3e5bd07e8c929730040b60d16791b6c0616d00
+21 -18
View File
@@ -155,11 +155,18 @@ threads-required = "num-test-threads"
filter = 'package(rustfs) & test(=connect::diagnostics::trace_runtime::tests::local_runtime_profile_cancellation_waits_for_lease_release)'
threads-required = "num-test-threads"
# These concurrent state writers keep the production 5s lock-acquisition limit
# while their peer completes durable IO. Reserve capacity from unrelated test
# processes while preserving each test's internal two-writer race.
# These real IAM/state fixtures keep production lock-acquisition limits while
# peers complete durable IO. Reserve capacity from unrelated test processes;
# the shared fixture's durable_iam_state_ prefix keeps the whole family together.
[[profile.default.overrides]]
filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::test_(peer_edit_generations_are_unique_across_nodes|recreated_state_object_allocates_over_the_previous_lifetimes_mark|state_object_lock_serializes_writers_from_separate_nodes|retry_event_persist_must_not_wipe_concurrent_locked_rmw)$/)'
filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::durable_iam_state_/)'
threads-required = "num-test-threads"
# Real target-repair and rebalance-stop fixtures commit durable metadata while
# peers wait on production locks or bounded cancellation checks. Reserve IO
# capacity for each fixture; its internal races and deadlines remain intact.
[[profile.default.overrides]]
filter = 'package(rustfs) & test(/^admin::handlers::(replication::target_repair_tests::|rebalance::rebalance_handler_tests::real_admin_stop_)/)'
threads-required = "num-test-threads"
# Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129).
@@ -319,13 +326,6 @@ retries = 2
filter = 'package(rustfs-ecstore) & (test(concurrent_resend_same_part_commits_one_generation) | test(concurrent_config_writes_from_separate_nodes_do_not_lose_writes))'
test-group = 'ecstore-serial-flaky'
# QUARANTINE: OPEN rustfs#4690 — walk_dir stall-budget accounting test depends
# on producer/consumer timing windows that stretch past the budget on loaded
# CI runners (regression test for rustfs#4644; failed on a zero-Rust-diff PR).
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & test(walk_dir_does_not_charge_consumer_backpressure_to_the_stall_budget)'
retries = 2
# Serialize the relocated-pool GET resume regression under the ci profile too
# (see the matching default-profile override near the top). No longer a
# quarantine: the fixture race (rustfs#6701/rustfs#6703) was fixed by #6707,
@@ -384,10 +384,16 @@ threads-required = "num-test-threads"
filter = 'package(rustfs) & test(=connect::diagnostics::trace_runtime::tests::local_runtime_profile_cancellation_waits_for_lease_release)'
threads-required = "num-test-threads"
# Keep the same state-writer capacity reservation in CI without changing the
# Keep the same IAM/state fixture capacity reservation in CI without changing the
# production lock deadline, internal concurrency, assertions, or retry policy.
[[profile.ci.overrides]]
filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::test_(peer_edit_generations_are_unique_across_nodes|recreated_state_object_allocates_over_the_previous_lifetimes_mark|state_object_lock_serializes_writers_from_separate_nodes|retry_event_persist_must_not_wipe_concurrent_locked_rmw)$/)'
filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::durable_iam_state_/)'
threads-required = "num-test-threads"
# Match the real admin fixture capacity reservation without retries or changes
# to production lock deadlines, cancellation barriers, or assertions.
[[profile.ci.overrides]]
filter = 'package(rustfs) & test(/^admin::handlers::(replication::target_repair_tests::|rebalance::rebalance_handler_tests::real_admin_stop_)/)'
threads-required = "num-test-threads"
# Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129).
@@ -536,13 +542,10 @@ fail-fast = false
[profile.e2e-smoke.junit]
path = "junit.xml"
# The pagination boundary cases can stall when a server/listing regression
# prevents the continuation request from completing. Keep the timeout scoped
# to those known failure modes so legitimate lifecycle/tiering waits retain
# their test-level timing budget.
# These storage-heavy fixtures reserve capacity. Pagination tests bound listing
# requests separately, after their seed objects have been written.
[[profile.e2e-smoke.overrides]]
filter = 'package(e2e_test) & test(/^list_objects_v2_pagination_test::tests::(test_list_objects_v2_delimiter_small_page_traverses_all|test_list_objects_v2_max_keys_above_limit_returns_token|test_list_objects_v2_maxkeys_above_limit_with_delimiter)$/)'
slow-timeout = { period = "60s", terminate-after = 2, grace-period = "10s" }
threads-required = "num-test-threads"
[[profile.e2e-smoke.overrides]]
+26 -15
View File
@@ -193,13 +193,10 @@ jobs:
# Raised to 3 on main pushes and manual dispatches; PRs keep 2 so the
# merge path is untouched while the experiment runs.
#
# Dispatch is included because push alone cannot supply the samples:
# this workflow cancels superseded runs on main, and only 4 of the last
# 20 push-triggered Test and Lint jobs reached a terminal state — at
# that rate ten samples would take roughly fifty merges. The
# concurrency group is scoped by event_name, so a dispatched run has
# its own group and is not cancelled by merge traffic, which makes the
# sample collectable on demand rather than by waiting.
# During baseline collection, main pushes cancelled superseded runs:
# only 4 of 20 Test and Lint jobs reached a terminal state. Running
# main validations now finish; dispatch remains an independent group
# for collecting additional samples on demand.
#
# Baseline over 17 samples at 2:
# median nextest/clippy step ratio 1.95, spread 1.85-2.06. The gate-2
@@ -248,11 +245,11 @@ jobs:
mkdir -p artifacts/test-and-lint
set +e
timeout --verbose --signal=TERM --kill-after=30s 15m \
cargo test --all --doc \
cargo test --all --exclude e2e_test --doc \
2>&1 | tee artifacts/test-and-lint/doctest.log
status=${PIPESTATUS[0]}
{
echo "command=cargo test --all --doc"
echo "command=cargo test --all --exclude e2e_test --doc"
echo "exit_status=${status}"
echo "finished_at=$(date --utc --iso-8601=seconds)"
echo
@@ -537,11 +534,15 @@ jobs:
cache-save-if: 'false'
install-build-packaging-tools: 'false'
- name: Protect Connect test home
run: chmod go-w "$(realpath "$HOME")"
- name: Prepare protocol test state
run: |
rm -f target/nextest/ci/junit.xml
chmod go-w "$(realpath "$HOME")"
- name: Run clippy with ${{ matrix.features.name }}
run: |
./scripts/ci/resource_sampler.sh start clippy
trap './scripts/ci/resource_sampler.sh stop' EXIT
cargo clippy -p rustfs -p rustfs-protocols --all-targets ${{ matrix.features.flags }} -- -D warnings
- name: Run tests with ${{ matrix.features.name }}
@@ -551,24 +552,26 @@ jobs:
# Keep feature-test linking under the same bounded concurrency as the
# main nextest lane; Clippy is metadata-only and needs no such limit.
CARGO_BUILD_JOBS: "2"
SAMPLER_INTERVAL_SECS: "15"
run: |
# --profile ci so the quarantine list (and its junit flaky markers)
# covers this leg too; the default profile is the local no-retry
# profile and silently ignored quarantined flakes here (rustfs#6703).
mkdir -p artifacts/protocol-tests
rm -f target/nextest/ci/junit.xml
./scripts/ci/resource_sampler.sh start nextest
trap './scripts/ci/resource_sampler.sh stop' EXIT
cargo nextest run --profile ci -p rustfs -p rustfs-protocols ${{ matrix.features.flags }} \
2>&1 | tee artifacts/protocol-tests/nextest.log
- name: Upload protocol test reports and diagnostics
if: >-
always() && contains(fromJSON('["success", "failure", "cancelled"]'), steps.protocol-tests.outcome)
if: always()
uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
with:
name: junit-test-and-lint-${{ matrix.features.name }}-${{ github.run_number }}-${{ github.run_attempt }}
path: |
target/nextest/ci/junit.xml
artifacts/protocol-tests
artifacts/test-and-lint
retention-days: 3
if-no-files-found: warn
@@ -827,6 +830,10 @@ jobs:
python3 ./scripts/check_test_wiring.py --check-profile e2e-smoke "${NEXTEST_LISTING}"
./scripts/check_security_smoke_count.sh check "${NEXTEST_LISTING}"
# Keep E2E documentation coverage with its already-built package graph.
- name: Run E2E documentation tests
run: cargo test -p e2e_test --doc
# PR smoke subset of the in-repo e2e suite (backlog#1149 ci-4). The
# profile.e2e-smoke default-filter in .config/nextest.toml is the single
# wiring mechanism for e2e tests in CI — extend that filter instead of
@@ -836,18 +843,22 @@ jobs:
env:
NEXTEST_ARCHIVE: ${{ runner.temp }}/rustfs-e2e-smoke.tar.zst
RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-smoke-logs
SAMPLER_INTERVAL_SECS: "15"
run: |
./scripts/ci/resource_sampler.sh start e2e-smoke
trap './scripts/ci/resource_sampler.sh stop' EXIT
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-smoke --archive-file "${NEXTEST_ARCHIVE}" \
--status-level all --final-status-level all --failure-output final
- name: Upload e2e smoke diagnostics
if: failure()
if: always()
uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
with:
name: e2e-smoke-diagnostics-${{ github.run_number }}
path: |
${{ runner.temp }}/rustfs-e2e-smoke-logs/
${{ runner.temp }}/rustfs-e2e-smoke-list.json
artifacts/test-and-lint
if-no-files-found: warn
- name: Upload e2e smoke JUnit report
@@ -31,8 +31,19 @@ mod tests {
use aws_sdk_s3::Client;
use aws_sdk_s3::primitives::ByteStream;
use std::collections::HashSet;
use std::future::Future;
use std::time::Duration;
use tokio::time::{Instant, error::Elapsed, timeout_at};
use tracing::info;
const PAGINATION_TIMEOUT: Duration = Duration::from_secs(120);
// Fixture population can require thousands of durable PUTs. Bound the
// listing traversal itself with one deadline shared by every page.
async fn listing_request_before<T>(deadline: Instant, request: impl Future<Output = T>) -> Result<T, Elapsed> {
timeout_at(deadline, request).await
}
/// Helper function to create an S3 client for testing
fn create_s3_client(env: &RustFSTestEnvironment) -> Client {
env.create_s3_client()
@@ -57,6 +68,54 @@ mod tests {
}
}
#[tokio::test]
async fn test_list_objects_v2_deadline_rejects_unresponsive_peer() {
use tokio::io::AsyncReadExt;
use tokio::net::TcpListener;
use tokio::sync::oneshot;
use tokio::time::timeout;
let listener = TcpListener::bind("127.0.0.1:0")
.await
.expect("bind unresponsive listing peer");
let endpoint = format!("http://{}", listener.local_addr().expect("listing peer address"));
let (received_tx, received_rx) = oneshot::channel();
let server = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.expect("accept listing request");
let mut bytes = [0; 4096];
assert!(stream.read(&mut bytes).await.expect("read listing request") > 0);
received_tx.send(()).expect("signal received listing request");
std::future::pending::<()>().await;
drop(stream);
});
let client = Client::from_conf(
crate::common::build_test_s3_config(&endpoint, "test-access", "test-secret", None, "pagination-deadline")
.to_builder()
.retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1))
.build(),
);
let request = client.list_objects_v2().bucket("deadline-test").send();
tokio::pin!(request);
tokio::select! {
response = &mut request => panic!("unresponsive listing peer unexpectedly returned: {response:?}"),
received = timeout(Duration::from_secs(10), received_rx) => {
received.expect("listing request must reach peer").expect("listing peer must remain alive");
}
}
// Expire the same absolute deadline after the real request is observed;
// this tests a stalled response without relying on transport timing.
let result = timeout(Duration::from_secs(1), listing_request_before(Instant::now(), request))
.await
.expect("expired pagination deadline must stop the in-flight request");
server.abort();
let _ = server.await;
assert!(
result.is_err(),
"an unresponsive ListObjectsV2 request must exceed the pagination deadline"
);
}
/// Test for Issue #2775: continuation forwarding must not
/// skip a child directory when the prefix component repeats in the key.
#[tokio::test]
@@ -514,12 +573,12 @@ mod tests {
.expect("Failed to put object");
}
let output = client
.list_objects_v2()
.bucket(bucket)
.max_keys(1001)
.send()
eprintln!("Seeded {object_count} objects in {bucket}; starting ListObjectsV2 pagination");
let deadline = Instant::now() + PAGINATION_TIMEOUT;
let output = listing_request_before(deadline, client.list_objects_v2().bucket(bucket).max_keys(1001).send())
.await
.expect("ListObjectsV2 pagination exceeded its deadline after fixture population")
.expect("Failed to list objects");
assert_eq!(output.contents().len(), 1000);
@@ -535,14 +594,18 @@ mod tests {
.expect("NextContinuationToken should be present when capped response is truncated")
.to_string();
let output = client
.list_objects_v2()
.bucket(bucket)
.max_keys(1001)
.continuation_token(next_token)
.send()
.await
.expect("Failed to list objects with continuation token");
let output = listing_request_before(
deadline,
client
.list_objects_v2()
.bucket(bucket)
.max_keys(1001)
.continuation_token(next_token)
.send(),
)
.await
.expect("ListObjectsV2 continuation exceeded the shared pagination deadline")
.expect("Failed to list objects with continuation token");
assert_eq!(output.contents().len(), 2);
assert!(!output.is_truncated().unwrap_or(false));
@@ -765,6 +828,9 @@ mod tests {
}
}
eprintln!("Seeded {} objects in {bucket}; starting ListObjectsV2 pagination", all_keys.len());
let deadline = Instant::now() + PAGINATION_TIMEOUT;
// Paginate with delimiter="/" and max_keys=50
// Visible per page: up to 50 CommonPrefixes
let mut listed_keys = Vec::new();
@@ -780,7 +846,10 @@ mod tests {
request = request.continuation_token(token);
}
let output = request.send().await.expect("Failed to list objects");
let output = listing_request_before(deadline, request.send())
.await
.expect("ListObjectsV2 delimiter traversal exceeded the shared pagination deadline")
.expect("Failed to list objects");
last_page_is_truncated = output.is_truncated().unwrap_or(false);
for obj in output.contents() {
@@ -988,15 +1057,18 @@ mod tests {
}
}
eprintln!(
"Seeded {} objects in {bucket}; starting ListObjectsV2 pagination",
dir_count * files_per_dir
);
let deadline = Instant::now() + PAGINATION_TIMEOUT;
// With delimiter: 12 CommonPrefixes visible, all fit within capped 1000
let output = client
.list_objects_v2()
.bucket(bucket)
.delimiter("/")
.max_keys(2000)
.send()
.await
.expect("Failed to list objects");
let output =
listing_request_before(deadline, client.list_objects_v2().bucket(bucket).delimiter("/").max_keys(2000).send())
.await
.expect("ListObjectsV2 delimiter listing exceeded its deadline after fixture population")
.expect("Failed to list objects");
assert_eq!(
output.common_prefixes().len(),
+11 -1
View File
@@ -1056,7 +1056,17 @@ async fn test_hermetic_transition_restore_failure_expiry_and_retry() -> TestResu
cold_client.create_bucket().bucket(TIER_BUCKET).send().await?;
let mut hot = RustFSTestEnvironment::new().await?;
start_tier_source(&mut hot, &[("RUSTFS_SCANNER_CYCLE", "1"), ("RUSTFS_ILM_DEBUG_DAY_SECS", "5")]).await?;
// Revisit the restored key every cycle; the default 16-cycle sampling is
// independent of the accelerated lifecycle clock used by this fixture.
start_tier_source(
&mut hot,
&[
("RUSTFS_SCANNER_CYCLE", "1"),
("RUSTFS_DATA_USAGE_UPDATE_DIR_CYCLES", "1"),
("RUSTFS_ILM_DEBUG_DAY_SECS", "5"),
],
)
.await?;
let hot_client = hot.create_s3_client();
add_rustfs_tier(&hot, &cold).await?;
@@ -51,6 +51,7 @@ const PLAIN_BUCKET: &str = "upgrade-plain-data";
const VERSIONED_BUCKET: &str = "upgrade-versioned-data";
const MIXED_BUCKET: &str = "upgrade-mixed-version-data";
const MIXED_NODE_COUNT: usize = 4;
const UPGRADE_READINESS_BODY: &[u8] = b"upgrade write readiness";
const MULTIPART_WORKERS: usize = 16;
const MULTIPART_UPLOADS_PER_WORKER: usize = 16;
// Peers keep a restarted node's drive in Suspect/Returning for roughly
@@ -442,6 +443,16 @@ async fn exercise_mixed_cluster(
Ok(())
}
fn upgrade_probe_client(client: &Client) -> Client {
Client::from_conf(
client
.config()
.to_builder()
.retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1))
.build(),
)
}
async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budget: Duration) -> TestResult {
// ListBuckets can succeed before peers recover a restarted disk. Cluster
// health also accepts Returning disks whose write health is still FAULTY.
@@ -449,13 +460,7 @@ async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budge
// writes still execute once and retain their original assertions.
let deadline = Instant::now() + budget;
for (node, client) in clients.iter().enumerate() {
let client = Client::from_conf(
client
.config()
.to_builder()
.retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1))
.build(),
);
let client = upgrade_probe_client(client);
let key = format!(".upgrade-readiness/{phase}/node-{node}");
let mut last_response = "no response".to_string();
loop {
@@ -470,7 +475,7 @@ async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budge
.put_object()
.bucket(MIXED_BUCKET)
.key(&key)
.body(ByteStream::from_static(b"upgrade write readiness"))
.body(ByteStream::from_static(UPGRADE_READINESS_BODY))
.send(),
)
.await;
@@ -497,6 +502,45 @@ async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budge
Ok(())
}
async fn prepare_previous_release_baseline(cluster: &mut RustFSTestClusterEnvironment, previous_binary: &Path) -> TestResult {
let clients: Vec<_> = cluster.create_all_clients()?.iter().map(upgrade_probe_client).collect();
// rc.5 can latch a non-elected node's write fence while its first startup
// waits for pool metadata (#7473). First prove the elected writer can
// persist data, then restart each old process once with three peers still
// readable. This preparation ends before any current binary is started.
wait_for_upgrade_write_readiness(&clients[..1], "previous-seed", LISTING_CONVERGENCE_TIMEOUT).await?;
for node in [1, 2, 3, 0] {
cluster.stop_node(node)?;
cluster.start_node_from_binary(node, previous_binary).await?;
wait_for_upgrade_write_readiness(
std::slice::from_ref(&clients[node]),
&format!("previous-restart-{node}"),
LISTING_CONVERGENCE_TIMEOUT,
)
.await?;
}
wait_for_upgrade_write_readiness(&clients, "previous-baseline", LISTING_CONVERGENCE_TIMEOUT).await?;
for (reader, client) in clients.iter().enumerate() {
assert_eq!(
read_object(client, MIXED_BUCKET, ".upgrade-readiness/previous-seed/node-0", None)
.await?
.1,
UPGRADE_READINESS_BODY,
"previous-release node {reader} must retain the seed across its rolling restart"
);
for writer in 0..clients.len() {
let key = format!(".upgrade-readiness/previous-baseline/node-{writer}");
assert_eq!(
read_object(client, MIXED_BUCKET, &key, None).await?.1,
UPGRADE_READINESS_BODY,
"previous-release node {reader} must read baseline data from writer {writer}"
);
}
}
Ok(())
}
#[cfg(test)]
mod upgrade_write_readiness_tests {
use super::*;
@@ -875,6 +919,7 @@ async fn rolling_upgrade_from_rc2_preserves_mixed_version_contracts() -> TestRes
configure_cluster_logs(&mut cluster)?;
cluster.start_with_binary(&previous_binary).await?;
cluster.create_test_bucket(MIXED_BUCKET).await?;
prepare_previous_release_baseline(&mut cluster, &previous_binary).await?;
cluster.stop_node(0)?;
cluster.start_node_from_binary(0, &current_binary).await?;
+20 -10
View File
@@ -13428,24 +13428,24 @@ mod test {
);
}
/// A writer that stalls on every write, standing in for a slow listing
/// consumer (quorum merge, a lagging peer drive).
/// A writer that advances paused time on every write, standing in for a
/// slow listing consumer without depending on host scheduling.
struct SlowWriter {
delay: Duration,
sleep: Option<Pin<Box<Sleep>>>,
advance: Option<Pin<Box<dyn std::future::Future<Output = ()> + Send>>>,
}
impl AsyncWrite for SlowWriter {
fn poll_write(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<io::Result<usize>> {
if self.sleep.is_none() {
if self.advance.is_none() {
let delay = self.delay;
self.sleep = Some(Box::pin(tokio::time::sleep(delay)));
self.advance = Some(Box::pin(tokio::time::advance(delay)));
}
let sleep = self.sleep.as_mut().expect("sleep was just installed");
match sleep.as_mut().poll(cx) {
let advance = self.advance.as_mut().expect("clock advance was just installed");
match advance.as_mut().poll(cx) {
Poll::Ready(()) => {
self.sleep = None;
self.advance = None;
Poll::Ready(Ok(buf.len()))
}
Poll::Pending => Poll::Pending,
@@ -13480,6 +13480,7 @@ mod test {
let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse");
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
disk.wait_for_startup_cleanup().await;
let stall = Duration::from_millis(300);
let write_delay = Duration::from_millis(150);
@@ -13493,12 +13494,21 @@ mod test {
let mut writer = SlowWriter {
delay: write_delay,
sleep: None,
advance: None,
};
let started = std::time::Instant::now();
// A live blocking task inhibits Tokio's automatic clock advance while
// real filesystem reads are pending. Only consumer writes advance time.
let (clock_guard_tx, clock_guard_rx) = std::sync::mpsc::channel::<()>();
let clock_guard = tokio::task::spawn_blocking(move || clock_guard_rx.recv());
tokio::time::pause();
let started = Instant::now();
let result = disk.walk_dir(opts, &mut writer).await;
let elapsed = started.elapsed();
drop(clock_guard_tx);
let _ = clock_guard
.await
.expect("clock guard should exit after its sender is dropped");
assert!(result.is_ok(), "a walk making steady progress must not time out, got {result:?}");
assert!(
@@ -378,8 +378,10 @@ async fn blackbox_heal_requests_preserve_repair_scope() {
// admission channel and wait for its receipt before acknowledging the write.
// Drive the test receiver alongside the PUT so neither waits on the other.
let (_put_dirs, put_store) = crate::bucket::metadata_sys::test_support::isolated_store_over_temp_disks().await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(put_store.clone(), Vec::new()).await;
let put_set = put_store.pools[0].disk_set[0].clone();
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&put_store), Vec::new()).await;
let put_set = Arc::clone(&put_store.pools[0].disk_set[0]);
assert_eq!(put_set.set_drive_count, 4);
assert_eq!(put_set.default_parity_count, 2);
let put_bucket = "bb-put-partial-convergence";
let put_object = "object.bin";
put_store
@@ -389,7 +391,7 @@ async fn blackbox_heal_requests_preserve_repair_scope() {
let put_incarnation = put_store
.bucket_incarnation_id_from_disk(put_bucket)
.await
.expect("PUT bucket should have a persisted incarnation");
.expect("partial repair must bind the real persisted bucket generation");
let offline_disk = {
let mut disks = put_set.disks.write().await;
disks[0].take()
@@ -423,9 +425,10 @@ async fn blackbox_heal_requests_preserve_repair_scope() {
.to_string();
assert_eq!(request.object_version_id.as_deref(), Some(committed_version.as_str()));
assert_eq!(request.expected_bucket_incarnation_id, Some(put_incarnation));
assert_eq!(request.pool_index, Some(0));
assert_eq!(request.set_index, Some(0));
assert_eq!(request.source, HealRequestSource::Mrf);
assert_eq!(request.expected_bucket_incarnation_id, Some(put_incarnation));
let duplicate_request = tokio::time::timeout(std::time::Duration::from_millis(100), async {
loop {
+69 -3
View File
@@ -3502,7 +3502,10 @@ impl SetDisks {
let protect_write = opts.shard_integrity_write_enabled();
let source_bucket_incarnation_id = match opts.expected_bucket_incarnation_id {
Some(incarnation_id) => Some(incarnation_id),
None => self.bucket_incarnation_id_from_disk(bucket).await.ok(),
// Internal metadata writes may already hold the pool metadata
// lock; resolving a user bucket identity would re-enter it.
None if !is_meta_bucketname(bucket) => self.bucket_incarnation_id_from_disk(bucket).await.ok(),
None => None,
};
if publication_fence.is_none()
&& opts.data_movement
@@ -8077,7 +8080,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
#[tracing::instrument(skip(self))]
async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> {
let source_bucket_incarnation_id = self.bucket_incarnation_id_from_disk(bucket).await.ok();
let source_bucket_incarnation_id = if is_meta_bucketname(bucket) {
None
} else {
self.bucket_incarnation_id_from_disk(bucket).await.ok()
};
self.delete_object_version_with_purge(bucket, object, fi, force_del_marker, None, source_bucket_incarnation_id)
.await
}
@@ -9277,7 +9284,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
#[tracing::instrument(skip(self))]
async fn add_partial(&self, bucket: &str, object: &str, version_id: &str) -> Result<()> {
let source_bucket_incarnation_id = self.bucket_incarnation_id_from_disk(bucket).await.ok();
let source_bucket_incarnation_id = if is_meta_bucketname(bucket) {
None
} else {
self.bucket_incarnation_id_from_disk(bucket).await.ok()
};
self.add_partial_with_source_incarnation(bucket, object, version_id, source_bucket_incarnation_id)
.await
}
@@ -22443,6 +22454,61 @@ mod single_delete_namespace_owner_tests {
.expect("seed a complete real object version");
}
#[tokio::test]
#[serial_test::serial]
async fn internal_metadata_mutations_do_not_reenter_pool_metadata() {
let (_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let set = &store.pools[0].disk_set[0];
let disks = set.disk_inventory().await.into_iter().flatten().collect::<Vec<_>>();
for bucket in [RUSTFS_META_BUCKET, crate::disk::MIGRATING_META_BUCKET] {
for disk in &disks {
if let Err(err) = disk.make_volume(bucket).await {
assert_eq!(err, DiskError::VolumeExists, "prepare the internal metadata volume");
}
}
let version = Uuid::new_v4();
seed_version(set, bucket, "delete-under-pool-lock", version, b"metadata version").await;
let request = FileInfo {
name: "delete-under-pool-lock".to_string(),
version_id: Some(version),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
};
let mut reader = PutObjReader::from_vec(b"metadata under lock".to_vec());
let opts = ObjectOptions {
write_completion: WriteCompletion::TailDrained,
..Default::default()
};
// Pool metadata transactions retain this write guard while saving
// internal objects. None of these paths may look up a user bucket.
let _pool_meta_guard = store.pool_meta.write().await;
let limit = Duration::from_secs(30);
let (put, delete, partial) = tokio::join!(
tokio::time::timeout(limit, set.put_object(bucket, "put-under-pool-lock", &mut reader, &opts)),
tokio::time::timeout(limit, set.delete_object_version(bucket, &request.name, &request, false)),
tokio::time::timeout(limit, set.add_partial(bucket, "partial-under-pool-lock", "")),
);
assert!(
matches!(&put, Ok(Ok(_))) && matches!(&delete, Ok(Ok(()))) && matches!(&partial, Ok(Ok(()))),
"internal operations must finish while pool metadata is locked: bucket={bucket}, put={put:?}, delete={delete:?}, partial={partial:?}"
);
for disk in &disks {
let written = disk
.read_version("", bucket, "put-under-pool-lock", "", &ReadOptions::default())
.await
.expect("the internal PUT must reach every disk");
assert_eq!(written.size, b"metadata under lock".len() as i64);
let deleted = disk
.read_version("", bucket, &request.name, &version.to_string(), &ReadOptions::default())
.await;
assert!(matches!(deleted, Err(DiskError::FileNotFound | DiskError::FileVersionNotFound)));
}
}
}
#[tokio::test]
#[serial_test::serial(capacity_dirty_scope)]
async fn single_delete_advances_namespace_generation_through_cleanup() {
+30 -1
View File
@@ -1872,8 +1872,12 @@ mod tests {
assert_eq!(residue.diagnostic_bytes_read, 0);
}
#[tokio::test]
#[tokio::test(start_paused = true)]
async fn bucket_residue_scan_distinguishes_visible_and_tier_free_xlmeta() {
// Classification keeps the production budget, but real filesystem
// scheduling must not advance the clock before classification finishes.
let (clock_guard_tx, clock_guard_rx) = std::sync::mpsc::channel::<()>();
let clock_guard = tokio::task::spawn_blocking(move || clock_guard_rx.recv());
let root = tempfile::tempdir().expect("temporary bucket root should be created");
let bucket_path = root.path().join("bucket");
let visible_path = bucket_path.join("visible").join(STORAGE_FORMAT_FILE);
@@ -1895,6 +1899,7 @@ mod tests {
let visible_scan = scan_metadata_less_residue(&bucket_path)
.await
.expect("visible xl.meta scan should succeed");
assert!(!visible_scan.diagnostic_truncated, "{visible_scan:?}");
assert_eq!(visible_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::VisibleVersion));
tokio::fs::remove_dir_all(bucket_path.join("visible"))
@@ -1931,6 +1936,7 @@ mod tests {
let free_scan = scan_metadata_less_residue(&bucket_path)
.await
.expect("free-version xl.meta scan should succeed");
assert!(!free_scan.diagnostic_truncated, "{free_scan:?}");
assert_eq!(free_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::TierFreeVersion));
tokio::fs::remove_dir_all(bucket_path.join("free"))
@@ -1951,6 +1957,7 @@ mod tests {
let exact_limit_scan = scan_metadata_less_residue(&bucket_path)
.await
.expect("exact-limit xl.meta scan should remain fail closed");
assert!(!exact_limit_scan.diagnostic_truncated, "{exact_limit_scan:?}");
assert_eq!(exact_limit_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::UnknownXlMeta));
assert_eq!(exact_limit_scan.diagnostic_bytes_read, BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES);
tokio::fs::remove_dir_all(bucket_path.join("exact-limit"))
@@ -1971,9 +1978,31 @@ mod tests {
let oversized_scan = scan_metadata_less_residue(&bucket_path)
.await
.expect("oversized xl.meta scan should remain fail closed");
assert!(!oversized_scan.diagnostic_truncated, "{oversized_scan:?}");
assert_eq!(oversized_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::UnknownXlMeta));
assert_eq!(oversized_scan.diagnostic_bytes_read, 0);
assert!(oversized_scan.diagnostic_bytes_read <= BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES);
drop(clock_guard_tx);
let _ = clock_guard
.await
.expect("classification clock guard should exit after its sender is dropped");
// An exhausted diagnostic budget may return before observing xl.meta.
// It must report truncation rather than inventing a classification.
let first_io_started = Arc::new(AtomicBool::new(false));
let mut budget = BucketDeleteDiagnosticBudget::new().with_first_io_delay(
BUCKET_DELETE_DIAGNOSTIC_MAX_ELAPSED + Duration::from_millis(100),
first_io_started.clone(),
);
let truncated = scan_metadata_less_residue_with_budget(&bucket_path, &mut budget)
.await
.expect("a diagnostic timeout should return a partial scan");
assert!(first_io_started.load(Ordering::SeqCst));
assert!(truncated.diagnostic_truncated, "{truncated:?}");
assert_eq!(truncated.xlmeta_blocker, None);
assert_eq!(truncated.entries_scanned, 0);
assert_eq!(truncated.diagnostic_bytes_read, 0);
}
#[tokio::test]
+7 -7
View File
@@ -3194,25 +3194,25 @@ mod tests {
replay_intent.kind = MrfKind::PartialWrite;
replay_intent.version_id = None;
let replay_payload = encoded_payload(&replay_intent);
let source_incarnation = storage
let source_bucket_incarnation_id = storage
.mrf_bucket_incarnation_id(bucket)
.await
.expect("read persisted replay source identity")
.expect("replay source bucket has a persisted incarnation");
.expect("read the producer's bucket incarnation")
.expect("the fixture bucket must have a persisted incarnation");
let lifecycle_limit = config.journal_max_bytes.saturating_mul(4);
let lifecycle = encode_mrf_lifecycle_checkpoint(
replay_owner,
11,
vec![ResponsibilityCheckpoint {
intent_digest: intent_digest(&replay_intent).expect("fixture intent has a canonical digest"),
intent_digest: intent_digest(&replay_intent).expect("encode the replay intent identity"),
responsibility_id: Uuid::new_v4(),
source_bucket_incarnation_id: Some(source_incarnation),
source_bucket_incarnation_id: Some(source_bucket_incarnation_id),
last_operator_acceptance: None,
state: partial_write::ResponsibilityState::Active,
}],
lifecycle_limit,
)
.expect("encode generation-bound replay responsibility");
.expect("encode the producer's source-bound responsibility");
snapshot::publish_committed_snapshot_with_companion(
&disks,
replay_owner,
@@ -3222,7 +3222,7 @@ mod tests {
Some((&MRF_LIFECYCLE_PATHS, &lifecycle, lifecycle_limit)),
)
.await
.expect("publish committed replay checkpoint with its source identity");
.expect("publish committed replay checkpoint with its producer identity");
let mut queue = MrfQueue::new(config.queue_capacity, config.journal_max_bytes);
let mut backoff_until = None;
@@ -513,14 +513,24 @@ async fn delete_during_root(case: DeleteCase) {
fn run_delete_case(case: DeleteCase) {
// Incarnation-bound repairs use a spawned owner; retain the server's stack
// budget while exercising the real EC2+2 encode/decode futures in debug.
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(4)
.thread_stack_size(8 * 1024 * 1024)
.enable_all()
.build()
.expect("concurrent delete runtime");
runtime.block_on(delete_during_root(case));
// budget for the whole scenario, including its foreground DELETE. block_on
// alone would poll that future on libtest's smaller calling-thread stack.
let stack_size = 8 * 1024 * 1024;
std::thread::Builder::new()
.name("heal-concurrent-delete".to_owned())
.stack_size(stack_size)
.spawn(move || {
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(4)
.thread_stack_size(stack_size)
.enable_all()
.build()
.expect("concurrent delete runtime");
runtime.block_on(delete_during_root(case));
})
.expect("concurrent delete scenario thread starts")
.join()
.expect("concurrent delete scenario must not panic");
}
#[test]
+121 -59
View File
@@ -229,42 +229,10 @@ fn committed_manifest(owner: uuid::Uuid, sequence: u64, payload: &[u8]) -> Vec<u
manifest
}
fn write_generation_bound_snapshot_to_disks(
disk_paths: &[PathBuf],
sequence: u64,
payload: &[u8],
source_incarnation: uuid::Uuid,
partial_records: &[&[u8]],
) {
assert!(!source_incarnation.is_nil(), "replay source must come from the persisted bucket");
let owner = uuid::Uuid::new_v4();
let records = partial_records
.iter()
.map(|record| {
serde_json::json!({
"intent_digest": Sha256::digest(record).to_vec(),
"responsibility_id": uuid::Uuid::new_v4(),
"source_bucket_incarnation_id": source_incarnation,
"last_operator_acceptance": null,
"state": { "state": "active" },
})
})
.collect::<Vec<_>>();
let lifecycle = serde_json::to_vec(&serde_json::json!({
"format_version": 1,
"checkpoint_owner": owner,
"checkpoint_sequence": sequence,
"records": records,
}))
.expect("encode source-bound lifecycle fixture");
let mut companion = b"RFMFLC01".to_vec();
companion.push(1);
companion.extend_from_slice(&u64::try_from(lifecycle.len()).expect("fixture length fits").to_le_bytes());
companion.extend_from_slice(&Sha256::digest(&lifecycle));
companion.extend_from_slice(&lifecycle);
write_journal_path_to_disks(disk_paths, ".heal-mrf-lifecycle.0.bin", &companion);
fn write_committed_snapshot_to_disks(disk_paths: &[std::path::PathBuf], sequence: u64, payload: &[u8]) {
let manifest = committed_manifest(uuid::Uuid::new_v4(), sequence, payload);
write_journal_path_to_disks(disk_paths, COMMITTED_PAYLOAD_REL, payload);
write_journal_path_to_disks(disk_paths, COMMITTED_MANIFEST_REL, &committed_manifest(owner, sequence, payload));
write_journal_path_to_disks(disk_paths, COMMITTED_MANIFEST_REL, &manifest);
}
fn journal_exists_on_all_disks(disk_paths: &[std::path::PathBuf], relative_path: &str) -> bool {
@@ -344,6 +312,87 @@ fn committed_checkpoint_matches_on_all_disks(disk_paths: &[PathBuf], sequence: u
})
}
async fn assert_legacy_responsibilities_parked(
disk_paths: &[PathBuf],
expected_records: &[&[u8]],
observed_bucket_incarnation_id: uuid::Uuid,
) {
assert!(!observed_bucket_incarnation_id.is_nil(), "the observed bucket generation must be real");
let checkpoint = mrf_queue::snapshot::inspect_local_committed_snapshot(rustfs_config::DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES)
.await
.expect("inspect the committed legacy responsibility checkpoint")
.expect("parked responsibilities must retain a committed checkpoint");
assert_eq!(
checkpoint.payload().len(),
expected_records.iter().map(|record| record.len()).sum::<usize>()
);
assert!(
expected_records
.iter()
.all(|record| { checkpoint.payload().windows(record.len()).any(|bytes| bytes == *record) }),
"the latest checkpoint must preserve each exact kind, scope, version and object identity"
);
// The production inspector selects and validates the committed slot. Read
// its private lifecycle companion without adding a test-only public API.
let companion_path = format!(".heal-mrf-lifecycle.{}.bin", checkpoint.slot());
let companions: Vec<_> = disk_paths
.iter()
.map(|path| std::fs::read(path.join(META_BUCKET).join(&companion_path)).expect("read each parked lifecycle replica"))
.collect();
let companion = companions.first().expect("the fixture must have journal disks");
assert!(
companions.iter().all(|replica| replica == companion),
"every disk must retain the same lifecycle responsibility"
);
assert!(companion.len() >= 49, "lifecycle envelope must be complete");
assert_eq!(&companion[..8], b"RFMFLC01");
assert_eq!(companion[8], 1);
assert_eq!(
u64::from_le_bytes(companion[9..17].try_into().expect("lifecycle payload length")),
u64::try_from(companion.len() - 49).expect("fixture lifecycle length fits")
);
assert_eq!(&companion[17..49], Sha256::digest(&companion[49..]).as_slice());
let payload: serde_json::Value = serde_json::from_slice(&companion[49..]).expect("decode committed lifecycle payload");
assert_eq!(payload["format_version"], 1);
assert_eq!(payload["checkpoint_owner"], checkpoint.owner().to_string());
assert_eq!(payload["checkpoint_sequence"], checkpoint.sequence());
let records = payload["records"].as_array().expect("lifecycle responsibility records");
assert_eq!(records.len(), expected_records.len());
let mut responsibility_ids = std::collections::HashSet::new();
for expected in expected_records {
let digest: [u8; 32] = Sha256::digest(expected).into();
let matching: Vec<_> = records
.iter()
.filter(|record| record["intent_digest"] == serde_json::json!(digest))
.collect();
assert_eq!(matching.len(), 1, "each exact partial-write identity must own one parked responsibility");
let record = matching[0];
let id = uuid::Uuid::parse_str(record["responsibility_id"].as_str().expect("responsibility UUID"))
.expect("valid responsibility UUID");
assert!(
!id.is_nil() && responsibility_ids.insert(id),
"scoped responsibilities must have distinct non-nil IDs"
);
assert_eq!(
record.get("source_bucket_incarnation_id"),
Some(&serde_json::Value::Null),
"replay must not invent the unknown producer generation"
);
assert_eq!(
record.get("last_operator_acceptance"),
Some(&serde_json::Value::Null),
"replay must not fabricate operator authorization"
);
assert_eq!(record["state"]["state"], "legacy_generation_unknown");
assert_eq!(
record["state"]["observed_bucket_incarnation_id"],
observed_bucket_incarnation_id.to_string()
);
assert!(record["state"]["detected_at_ms"].as_u64().is_some_and(|time| time > 0));
}
}
async fn wait_until<F, Fut>(deadline: Duration, mut probe: F) -> bool
where
F: FnMut() -> Fut,
@@ -393,12 +442,17 @@ async fn decode_failure_intent_maps_to_urgent_mrf_heal_request() {
/// A journal left behind by a previous process must be replayed into the
/// manager queue, and a torn tail must not block replay of the intact records.
/// The partial-write record keeps the legacy journal as the durable anchor
/// until an exact verified repair proof can discharge it.
/// A partial write with no producer generation remains durably parked rather
/// than being admitted against whichever bucket happens to exist at replay.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() {
let (disk_paths, storage) = heal_env_with_bucket("replay-bucket").await;
let bucket_incarnation_id = storage
.mrf_bucket_incarnation_id("replay-bucket")
.await
.expect("read the current replay bucket generation")
.expect("the replay bucket must have a persisted generation");
// The journal reader resolves disks through the process-local disk map;
// register the environment's disks the same way server startup does.
@@ -418,7 +472,10 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor()
assert_eq!(replayed, 2, "the two intact records must be replayed");
let snapshot = manager.operations_snapshot().await;
assert_eq!(snapshot.queued_by_source.mrf, 2, "replayed intents must be attributed to the MRF source");
assert_eq!(
snapshot.queued_by_source.mrf, 1,
"only the decode-failure intent may enter the MRF manager queue"
);
assert!(
disk_paths
@@ -441,7 +498,11 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor()
let snapshot = manager.operations_snapshot().await;
assert_eq!(snapshot.queued_by_priority.urgent, 1, "the decode-failure record must replay as Urgent");
assert!(snapshot.queued_by_priority.normal >= 1, "the partial-write record must replay as Normal");
assert_eq!(
snapshot.queued_by_priority.normal, 0,
"an unknown producer generation must not enter the repair queue"
);
assert_legacy_responsibilities_parked(&disk_paths, &[&successor], bucket_incarnation_id).await;
}
/// A committed checkpoint published by the new two-slot writer is the
@@ -451,16 +512,16 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor()
#[serial]
async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() {
let (disk_paths, storage) = heal_env_with_bucket("committed-bucket").await;
let bucket_incarnation_id = storage
.mrf_bucket_incarnation_id("committed-bucket")
.await
.expect("read the current committed bucket generation")
.expect("the committed bucket must have a persisted generation");
register_local_disks(&disk_paths, "mrf-committed-replay-test").await;
let committed = scoped_journal_record(3, "committed-bucket", "committed-object", Some([9u8; 16]), 0, 0, 0);
let stale_legacy = journal_record(1, "legacy-bucket", "legacy-object", None, 0);
let incarnation = storage
.mrf_bucket_incarnation_id("committed-bucket")
.await
.expect("read committed source incarnation")
.expect("committed source bucket has a persisted incarnation");
write_generation_bound_snapshot_to_disks(&disk_paths, 7, &committed, incarnation, &[&committed]);
write_committed_snapshot_to_disks(&disk_paths, 7, &committed);
write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &stale_legacy);
write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &stale_legacy);
@@ -469,10 +530,10 @@ async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() {
assert_eq!(replayed, 1, "only the committed snapshot epoch may replay");
let snapshot = manager.operations_snapshot().await;
assert_eq!(snapshot.queued_by_source.mrf, 1);
assert_eq!(snapshot.queued_by_source.mrf, 0);
assert_eq!(
snapshot.queued_by_priority.normal, 1,
"the committed partial-write record must replay instead of the stale legacy decode-failure"
snapshot.queued_by_priority.normal, 0,
"the committed partial write without a producer generation must remain parked"
);
assert_eq!(
snapshot.queued_by_priority.urgent, 0,
@@ -480,8 +541,9 @@ async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() {
);
assert!(
journal_exists_on_all_disks(&disk_paths, COMMITTED_MANIFEST_REL),
"the committed checkpoint remains until the accepted partial-write has proof"
"the committed checkpoint must retain the parked partial-write responsibility"
);
assert_legacy_responsibilities_parked(&disk_paths, &[&committed], bucket_incarnation_id).await;
}
/// A damaged committed checkpoint is ambiguous: replay must not fall back to
@@ -588,6 +650,11 @@ async fn authoritative_journal_is_not_merged_with_legacy_mirror() {
#[serial]
async fn authoritative_journal_replay_preserves_kind_and_scope_identity() {
let (disk_paths, storage) = heal_env_with_bucket("identity-bucket").await;
let bucket_incarnation_id = storage
.mrf_bucket_incarnation_id("identity-bucket")
.await
.expect("read the current identity bucket generation")
.expect("the identity bucket must have a persisted generation");
register_local_disks(&disk_paths, "mrf-authoritative-identity-test").await;
let first_partial = scoped_journal_record(3, "identity-bucket", "same-object", None, 0, 3, 7);
@@ -599,12 +666,6 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() {
let stale_legacy = journal_record(3, "identity-bucket", "stale-legacy-object", None, 0);
write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &authoritative);
write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &stale_legacy);
let incarnation = storage
.mrf_bucket_incarnation_id("identity-bucket")
.await
.expect("read authoritative source incarnation")
.expect("authoritative source bucket has a persisted incarnation");
write_generation_bound_snapshot_to_disks(&disk_paths, 1, &authoritative, incarnation, &[&first_partial, &second_partial]);
let manager = make_manager(storage);
let replayed = mrf_queue::replay_journal_once(&manager).await;
@@ -615,12 +676,12 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() {
let snapshot = manager.operations_snapshot().await;
assert_eq!(
snapshot.queued_by_source.mrf, 4,
"same-object MRF replay must retain distinct kind and scope responsibilities"
snapshot.queued_by_source.mrf, 2,
"only metadata-corruption and decode-failure repairs may enter the manager"
);
assert_eq!(
snapshot.queued_by_priority.normal, 2,
"the two scoped partial-write records must remain independently queued"
snapshot.queued_by_priority.normal, 0,
"neither scoped partial write may be repaired with an unknown producer generation"
);
assert_eq!(
snapshot.queued_by_priority.high, 1,
@@ -640,6 +701,7 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() {
journal_matches_on_all_disks(&disk_paths, JOURNAL_REL, &[]),
"the legacy mirror must not misrepresent scoped-only partial-write responsibilities"
);
assert_legacy_responsibilities_parked(&disk_paths, &[&first_partial, &second_partial], bucket_incarnation_id).await;
}
/// If replay reaches a full heal-manager queue, the old journal remains the
+43 -4
View File
@@ -2402,18 +2402,57 @@ mod tests {
assert!(pool_meta_lock.get_write_lock_quiet(Duration::from_millis(100)).await.is_err());
barrier.release();
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let loaded = loop {
let loaded = load_scanner_pause_backlog(Arc::clone(&store))
.await
.expect("cancellation cannot erase the old authority");
assert_eq!(loaded.ledger, original);
let committed = loaded.authoritative_commit.expect("seed retains the old cohort proof");
let committed = loaded
.authoritative_commit
.as_ref()
.expect("seed retains the old cohort proof");
assert_eq!(committed.replicas, vec![replica_id(0, 0), replica_id(0, 1)]);
assert_eq!(loaded.replica_count, 6);
if loaded.healthy_replicas == loaded.replica_count {
break;
break loaded;
}
tokio::task::yield_now().await;
}
};
let drained_fence = pool_meta_lock
.get_write_lock_quiet(Duration::from_secs(30))
.await
.expect("the detached seed owner releases its membership fence after persistence");
drop(drained_fence);
assert_eq!(loaded.ledger, original);
let committed = loaded
.authoritative_commit
.as_ref()
.expect("seed retains the old cohort proof");
assert_eq!(committed.replicas, vec![replica_id(0, 0), replica_id(0, 1)]);
assert_eq!(loaded.replica_count, 6);
assert_eq!(loaded.healthy_replicas, 6);
let seeded = loaded
.replicas
.iter()
.find(|replica| replica.id == replica_id(2, 0))
.expect("the canceled caller's admitted seed replica");
assert!(matches!(
&seeded.state,
ScannerPauseBacklogReplicaState::Valid(record)
if record.stable.as_ref() == Some(&original) && record.committed.as_ref() == Some(committed)
));
// The detached publication owner finishes every selected replica
// without advancing beyond the old committed authority.
let remaining = loaded
.replicas
.iter()
.find(|replica| replica.id == replica_id(2, 1))
.expect("the remaining seed replica");
assert!(matches!(
&remaining.state,
ScannerPauseBacklogReplicaState::Valid(record)
if record.stable.as_ref() == Some(&original) && record.committed.as_ref() == Some(committed)
));
})
.await
.expect("detached native seed owners must drain without the canceled caller");
@@ -16,7 +16,7 @@ use super::*;
use crate::storage_api::EcstoreHealResultItem as HealItem;
use crate::storage_api::scanner_io::BucketInfo;
use rustfs_common::mrf_channel::{
MrfIngressResult, MrfKind, MrfScope, note_mrf_repaired, take_mrf_repaired_events_for, try_send_mrf_intent_typed,
MrfScope, note_mrf_repaired, persist_partial_write_intent_with_incarnation, take_mrf_repaired_events_for,
};
use rustfs_heal::heal::{
manager::{HealConfig, HealManager},
@@ -379,30 +379,41 @@ async fn mrf_ownership_manager_completion_preserves_scanner_pending() {
1,
1,
));
let scope = Some(MrfScope {
let scope = MrfScope {
pool_index: 0,
set_index: 0,
});
assert_eq!(
try_send_mrf_intent_typed(MrfKind::PartialWrite, &bucket, object, Some(version), scope),
MrfIngressResult::Enqueued
);
// Re-admission establishes that the first terminal callback released
// its ingress lease. Statistics alone precede notice publication.
};
let before = manager.get_statistics().await;
let completed_before = before.successful_tasks + before.failed_tasks;
persist_partial_write_intent_with_incarnation(&bucket, object, Some(version), scope, Some(storage.bucket_incarnation_id))
.await
.expect("persist the partial write with its actual source incarnation");
// The scheduler publishes terminal repair notices before updating these
// counters. Earlier objects can only enter their blocked retry here.
tokio::time::timeout(Duration::from_secs(5), async {
loop {
match try_send_mrf_intent_typed(MrfKind::PartialWrite, &bucket, object, Some(version), scope) {
MrfIngressResult::Enqueued => break,
MrfIngressResult::Coalesced => tokio::task::yield_now().await,
other => panic!("unexpected retry ingress result: {other:?}"),
let called = storage
.calls
.lock()
.expect("fixture calls")
.get(*object)
.copied()
.unwrap_or(0);
let stats = manager.get_statistics().await;
if called > 0 && stats.successful_tasks + stats.failed_tasks > completed_before {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("production terminal releases its ingress lease");
.expect("production terminal publishes its repair notices and completion");
persist_partial_write_intent_with_incarnation(&bucket, object, Some(version), scope, Some(storage.bucket_incarnation_id))
.await
.expect("persist another generation for the same source incarnation");
tokio::time::timeout(Duration::from_secs(5), storage.retry_started.notified())
.await
.expect("the real consumer starts the second generation");
.expect("the real consumer starts the second object call");
assert!(
take_mrf_repaired_events_for(&bucket).is_empty(),
"{object}: task completion must not emit an unproved repair"
@@ -415,9 +426,9 @@ async fn mrf_ownership_manager_completion_preserves_scanner_pending() {
);
}
assert_eq!(
try_send_mrf_intent_typed(MrfKind::PartialWrite, &bucket, object, Some(version), scope),
MrfIngressResult::Coalesced,
"the in-flight retry retains its new ingress lease"
storage.calls.lock().expect("fixture calls").get(*object).copied(),
Some(2),
"the second object call remains blocked without a duplicate execution"
);
note_mrf_repaired(&bucket, object, Some(*version.as_bytes()));
let syncs_before_retry = scanner.pending_heal_sync_count;
+18 -16
View File
@@ -8878,6 +8878,8 @@ mod tests {
}
/// Publish a ready IAM app context so `apply_iam_item` gets past its IAM guard.
/// Tests using this real storage fixture keep the `durable_iam_state_` prefix so
/// nextest reserves capacity while their internal races retain production deadlines.
async fn publish_ready_iam_context() {
use crate::admin::runtime_sources::{AppContext, publish_test_app_context};
use rustfs_iam::store::{Store as _, object::IAM_CONFIG_PREFIX};
@@ -9077,7 +9079,7 @@ mod tests {
/// every gated item type.
#[tokio::test]
#[serial]
async fn apply_iam_item_applies_delayed_in_order_updates_through_the_receiver() {
async fn durable_iam_state_apply_iam_item_applies_delayed_in_order_updates_through_the_receiver() {
publish_ready_iam_context().await;
seed_two_peer_state_for_iam_apply().await;
let iam = current_iam_handle().expect("test IAM");
@@ -9158,7 +9160,7 @@ mod tests {
#[tokio::test]
#[serial]
async fn apply_group_member_add_preserves_a_disabled_group_status() {
async fn durable_iam_state_apply_group_member_add_preserves_a_disabled_group_status() {
publish_ready_iam_context().await;
seed_two_peer_state_for_iam_apply().await;
let iam = current_iam_handle().expect("test IAM");
@@ -9193,7 +9195,7 @@ mod tests {
/// membership-only guard above exists to prevent.
#[tokio::test]
#[serial]
async fn apply_group_snapshot_propagates_a_disabled_status_with_members() {
async fn durable_iam_state_apply_group_snapshot_propagates_a_disabled_status_with_members() {
publish_ready_iam_context().await;
seed_two_peer_state_for_iam_apply().await;
let iam = current_iam_handle().expect("test IAM");
@@ -9228,7 +9230,7 @@ mod tests {
/// the write of one cannot interleave with the other's.
#[tokio::test]
#[serial]
async fn apply_iam_item_serializes_a_concurrent_older_grant_and_newer_revoke() {
async fn durable_iam_state_apply_iam_item_serializes_a_concurrent_older_grant_and_newer_revoke() {
publish_ready_iam_context().await;
seed_two_peer_state_for_iam_apply().await;
let t1 = OffsetDateTime::now_utc() - time::Duration::hours(2);
@@ -9267,7 +9269,7 @@ mod tests {
/// applied the delete; a genuinely newer create still lands.
#[tokio::test]
#[serial]
async fn apply_iam_item_rejects_a_stale_recreate_after_a_replicated_delete() {
async fn durable_iam_state_apply_iam_item_rejects_a_stale_recreate_after_a_replicated_delete() {
publish_ready_iam_context().await;
seed_two_peer_state_for_iam_apply().await;
let iam = current_iam_handle().expect("test IAM");
@@ -9312,7 +9314,7 @@ mod tests {
/// then switched off, and carries the source stamp.
#[tokio::test]
#[serial]
async fn apply_iam_item_creates_a_replicated_service_account_with_its_status() {
async fn durable_iam_state_apply_iam_item_creates_a_replicated_service_account_with_its_status() {
publish_ready_iam_context().await;
seed_two_peer_state_for_iam_apply().await;
let iam = current_iam_handle().expect("test IAM");
@@ -9341,7 +9343,7 @@ mod tests {
#[tokio::test]
#[serial]
async fn apply_iam_item_accepts_minio_sts_account_item_type() {
async fn durable_iam_state_apply_iam_item_accepts_minio_sts_account_item_type() {
publish_ready_iam_context().await;
// MinIO madmin-go sends `SRIAMItemSTSAcc = "sts-account"`. The bogus session token
@@ -9363,7 +9365,7 @@ mod tests {
#[tokio::test]
#[serial]
async fn apply_iam_item_still_accepts_legacy_sts_credential_item_type() {
async fn durable_iam_state_apply_iam_item_still_accepts_legacy_sts_credential_item_type() {
publish_ready_iam_context().await;
// Older RustFS peers emit `sts-credential`; the alias stays accepted permanently
@@ -15400,7 +15402,7 @@ mod tests {
/// would delete the retry queue and every other field along with it.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial]
async fn test_missed_pending_clear_must_not_rewrite_the_state_object() {
async fn durable_iam_state_test_missed_pending_clear_must_not_rewrite_the_state_object() {
publish_ready_iam_context().await;
// One peer, no pending records: exactly the shape the persist helper's clear
@@ -15461,7 +15463,7 @@ mod tests {
/// write before A finishes.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn test_peer_join_admission_serializes_iam_apply_against_a_newer_join() {
async fn durable_iam_state_test_peer_join_admission_serializes_iam_apply_against_a_newer_join() {
publish_ready_iam_context().await;
// Whole-second timestamps so the RFC3339 round trip through the state
@@ -15565,7 +15567,7 @@ mod tests {
/// admission: the assertion on the empty IAM log turns red.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn test_peer_join_admission_serializes_across_separate_nodes() {
async fn durable_iam_state_test_peer_join_admission_serializes_across_separate_nodes() {
publish_ready_iam_context().await;
let now = OffsetDateTime::now_utc().replace_nanosecond(0).expect("truncate nanos");
@@ -15662,7 +15664,7 @@ mod tests {
/// acked rotation is cleared in the same transaction that reports true.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial]
async fn test_finalize_pending_rotation_three_way_contract() {
async fn durable_iam_state_test_finalize_pending_rotation_three_way_contract() {
publish_ready_iam_context().await;
let local_peer = PeerInfo {
@@ -15744,7 +15746,7 @@ mod tests {
/// update.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn test_state_object_lock_serializes_writers_from_separate_nodes() {
async fn durable_iam_state_test_state_object_lock_serializes_writers_from_separate_nodes() {
publish_ready_iam_context().await;
let seed = SiteReplicationState {
@@ -15796,7 +15798,7 @@ mod tests {
/// leave the receiver unable to tell which edit is newer.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn test_peer_edit_generations_are_unique_across_nodes() {
async fn durable_iam_state_test_peer_edit_generations_are_unique_across_nodes() {
publish_ready_iam_context().await;
// A configured site: `persist_site_replication_state_no_lock` clears
// the object once a site drops below two peers, and a cleared object
@@ -15847,7 +15849,7 @@ mod tests {
/// accepts the restarted counter instead of fencing it.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn test_recreated_state_object_allocates_over_the_previous_lifetimes_mark() {
async fn durable_iam_state_test_recreated_state_object_allocates_over_the_previous_lifetimes_mark() {
publish_ready_iam_context().await;
let seed = || SiteReplicationState {
peers: ["site-a", "site-b"]
@@ -15912,7 +15914,7 @@ mod tests {
/// (the red-light commit pinned the exact interleaving).
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn test_retry_event_persist_must_not_wipe_concurrent_locked_rmw() {
async fn durable_iam_state_test_retry_event_persist_must_not_wipe_concurrent_locked_rmw() {
publish_ready_iam_context().await;
const ROUNDS: usize = 8;
+53 -16
View File
@@ -41,6 +41,9 @@ use tokio::sync::watch;
use tokio_rustls::TlsAcceptor;
use tokio_util::sync::CancellationToken;
// Allow durable state and real TLS setup independently of transport and cancellation deadlines.
const STATUS_OBSERVATION_TIMEOUT: Duration = Duration::from_secs(10);
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";
@@ -191,9 +194,23 @@ impl Reply {
struct TestServer {
endpoint: String,
seen: Arc<Mutex<Vec<Value>>>,
paths: Arc<Mutex<Vec<String>>>,
task: tokio::task::JoinHandle<()>,
}
impl TestServer {
fn request_diagnostics(&self) -> String {
match self.paths.try_lock() {
Ok(paths) => format!(
"request_count={} request_paths={paths:?} server_task_finished={}",
paths.len(),
self.task.is_finished()
),
Err(error) => format!("request_paths_unavailable={error} server_task_finished={}", self.task.is_finished()),
}
}
}
impl Drop for TestServer {
fn drop(&mut self) {
self.task.abort();
@@ -207,16 +224,20 @@ async fn server(pki: &TestPki, replies: Vec<Reply>) -> TestServer {
let replies = Arc::new(Mutex::new(VecDeque::from(replies)));
let seen = Arc::new(Mutex::new(Vec::new()));
let captured = seen.clone();
let paths = Arc::new(Mutex::new(Vec::new()));
let captured_paths = paths.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();
let paths = captured_paths.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();
let paths = paths.clone();
async move {
let path = request.uri().path().to_owned();
assert!(
@@ -227,6 +248,7 @@ async fn server(pki: &TestPki, replies: Vec<Reply>) -> TestServer {
seen.lock()
.expect("seen lock")
.push(serde_json::from_slice(&body).expect("request JSON"));
paths.lock().expect("paths lock").push(path.clone());
let reply = if path.ends_with("/diagnosticExecutionReceipts") {
Reply {
status: StatusCode::OK,
@@ -266,6 +288,7 @@ async fn server(pki: &TestPki, replies: Vec<Reply>) -> TestServer {
TestServer {
endpoint: format!("https://localhost:{}/agent/", address.port()),
seen,
paths,
task,
}
}
@@ -314,23 +337,33 @@ fn summary() -> CoarseNodeSummary {
}
async fn wait_for(
server: &TestServer,
status: &mut watch::Receiver<HeartbeatStatus>,
predicate: impl Fn(&HeartbeatStatus) -> bool,
) -> HeartbeatStatus {
tokio::time::timeout(Duration::from_secs(3), async {
tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async {
loop {
let current = status.borrow_and_update().clone();
if predicate(&current) {
return current;
}
status
.changed()
.await
.unwrap_or_else(|error| panic!("heartbeat status channel closed: {error}; last status: {:?}", *status.borrow()));
status.changed().await.unwrap_or_else(|error| {
panic!(
"heartbeat status channel closed: {error}; last status: {:?}; {}",
*status.borrow(),
server.request_diagnostics()
)
});
}
})
.await
.unwrap_or_else(|error| panic!("heartbeat status timeout: {error}; last status: {:?}", *status.borrow()))
.unwrap_or_else(|error| {
panic!(
"heartbeat status timeout: {error}; last status: {:?}; {}",
*status.borrow(),
server.request_diagnostics()
)
})
}
async fn assert_credential_failure(config: HeartbeatConfig, server: &TestServer, expected: &str) {
@@ -340,7 +373,7 @@ async fn assert_credential_failure(config: HeartbeatConfig, server: &TestServer,
.expect("configured runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
wait_for(server, &mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
HeartbeatStatus::Failed { reason } if reason.contains(expected)
));
assert!(server.seen.lock().expect("seen lock").is_empty());
@@ -436,7 +469,7 @@ async fn corrupt_persisted_state_is_rejected_before_network_delivery() {
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
HeartbeatStatus::Failed { reason } if reason.contains("violates the protocol invariants")
));
assert!(server.seen.lock().expect("seen lock").is_empty());
@@ -505,7 +538,7 @@ async fn sends_only_l0_fields_and_accepts_additive_response_fields() {
let mut status = runtime.status();
assert_eq!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online {
server_time: "2038-01-19T03:14:07Z".to_owned()
}
@@ -599,7 +632,7 @@ async fn restart_replays_a_pending_heartbeat_from_before_environment_collection(
.expect("configured runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online { .. }
));
runtime.shutdown().await;
@@ -620,7 +653,7 @@ async fn restart_replays_pending_request_then_advances_sequence() {
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await;
wait_for(&first_server, &mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await;
runtime.shutdown().await;
let first = first_server.seen.lock().expect("seen lock")[0].clone();
drop(first_server);
@@ -730,7 +763,7 @@ async fn restart_replays_exact_legacy_capabilities_with_and_without_jobs() {
.expect("configured runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online { .. }
));
runtime.shutdown().await;
@@ -770,7 +803,7 @@ async fn legacy_compatibility_does_not_accept_changed_capabilities() {
let mut status = runtime.status();
assert!(
matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
HeartbeatStatus::Failed { reason } if reason.contains("violates the protocol invariants")
),
"accepted altered persisted capabilities: {mutation}, jobs={job_capable}, service_memory={service_memory}"
@@ -859,7 +892,7 @@ async fn retry_after_is_respected_with_the_local_upper_bound() {
.expect("configured runtime");
let mut status = runtime.status();
assert_eq!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await,
wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await,
HeartbeatStatus::BackingOff {
delay: Duration::from_millis(80)
}
@@ -887,7 +920,7 @@ async fn disconnects_use_exponential_backoff_with_a_cap() {
let mut status = runtime.status();
for delay in [20, 40, 80] {
assert_eq!(
wait_for(&mut status, |status| {
wait_for(&server, &mut status, |status| {
matches!(status, HeartbeatStatus::BackingOff { delay: observed } if *observed == Duration::from_millis(delay))
})
.await,
@@ -912,7 +945,11 @@ async fn revoked_credential_stops_and_exposes_local_status() {
.expect("configured runtime");
let mut status = runtime.status();
assert_eq!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::AuthenticationStopped { .. })).await,
wait_for(&server, &mut status, |status| matches!(
status,
HeartbeatStatus::AuthenticationStopped { .. }
))
.await,
HeartbeatStatus::AuthenticationStopped {
status: 401,
reason: Some("CREDENTIAL_REVOKED".to_owned())
+58 -3
View File
@@ -38,6 +38,9 @@ use zeroize::Zeroizing;
mod common;
static TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
// Real storage setup shares the measurement window. Verify functionality with
// an integration budget; deadline behavior is checked against a stalled peer.
const REAL_SERVER_MEASUREMENT_BUDGET: Duration = Duration::from_secs(10);
fn now() -> i64 {
SystemTime::now().duration_since(UNIX_EPOCH).expect("current time").as_secs() as i64
@@ -305,6 +308,56 @@ async fn deadline_and_in_flight_cancellation_stop_a_stalled_probe() {
assert_eq!(measurement.target.reason_code, ObjectTargetReasonCode::Cancelled);
}
#[tokio::test]
async fn s3_request_deadline_stops_a_peer_that_never_responds() {
use tokio::io::AsyncReadExt as _;
let _guard = TEST_LOCK.lock().await;
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("response listener");
let address = listener.local_addr().expect("listener address");
let (received, request_received) = tokio::sync::oneshot::channel();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.expect("object connection");
let mut request = [0_u8; 4096];
assert!(socket.read(&mut request).await.expect("request headers") > 0);
received.send(()).expect("report received request");
std::future::pending::<()>().await;
drop(socket);
});
let probe = S3ObjectProbe::new(
&format!("http://{address}"),
None,
None,
Zeroizing::new("access".to_owned()),
Zeroizing::new("secret".to_owned()),
Zeroizing::new(String::new()),
Duration::from_secs(60),
)
.expect("object probe");
let mut request = request(ObjectOperation::GetObject);
request.duration = Duration::from_secs(30);
let measurement = tokio::spawn(async move { measure_object(&request, &probe, &CancellationToken::new()).await });
// Complete real TCP setup before pausing time. The 30s operation deadline
// must fire independently of the client's longer 60s transport timeout.
tokio::time::timeout(Duration::from_secs(10), request_received)
.await
.expect("peer received the request")
.expect("peer notification");
tokio::time::pause();
tokio::time::advance(Duration::from_secs(30)).await;
let measurement = tokio::time::timeout(Duration::from_secs(1), measurement)
.await
.expect("operation deadline must stop the request")
.expect("measurement task")
.expect("typed timeout");
assert_eq!(measurement.result.outcome(), ObjectOutcome::Failed);
assert_eq!(measurement.target.reason_code, ObjectTargetReasonCode::TimedOut);
assert_eq!(measurement.target.completed_operations, 0);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn only_one_object_collector_can_run_at_a_time() {
let _guard = TEST_LOCK.lock().await;
@@ -519,6 +572,7 @@ fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put() {
async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_body() {
let _guard = TEST_LOCK.lock().await;
let startup = std::time::Instant::now();
let port = match find_available_port() {
Ok(port) => port,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
@@ -531,6 +585,7 @@ async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_bod
.build()
.await
.expect("start embedded server");
eprintln!("embedded object fixture ready after {:?}", startup.elapsed());
let probe = S3ObjectProbe::new(
&server.endpoint(),
None,
@@ -538,13 +593,13 @@ async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_bod
Zeroizing::new(server.access_key().to_owned()),
Zeroizing::new(server.secret_key().to_owned()),
Zeroizing::new(String::new()),
Duration::from_secs(2),
REAL_SERVER_MEASUREMENT_BUDGET,
)
.expect("object probe");
for operation in [ObjectOperation::GetObject, ObjectOperation::PutObject] {
let mut request = request(operation);
request.duration = Duration::from_secs(2);
request.duration = REAL_SERVER_MEASUREMENT_BUDGET;
if operation == ObjectOperation::PutObject {
request.artifact_uid = "019e3ae0-0000-7000-8000-000000000016".to_owned();
}
@@ -604,7 +659,7 @@ async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_bod
"--traffic-bytes",
"65536",
"--duration-millis",
"1000",
&REAL_SERVER_MEASUREMENT_BUDGET.as_millis().to_string(),
"--acknowledge-l1",
]);
let result = tokio::task::spawn_blocking(move || command.output())
+85 -14
View File
@@ -59,6 +59,9 @@ const EXPECTED_CERTIFICATE_LIFETIME_SECONDS: i64 = 86_400;
const EXPECTED_ROTATION_THRESHOLD_SECONDS: i64 = 120;
#[cfg(not(feature = "connect-e2e-short-credentials"))]
const EXPECTED_ROTATION_THRESHOLD_SECONDS: i64 = 8 * 60 * 60;
// Allow durable state and real TLS setup independently of transport and cancellation deadlines.
const STATUS_OBSERVATION_TIMEOUT: Duration = Duration::from_secs(10);
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";
@@ -209,6 +212,19 @@ struct TestServer {
task: tokio::task::JoinHandle<()>,
}
impl TestServer {
fn request_diagnostics(&self) -> String {
match self.paths.try_lock() {
Ok(paths) => format!(
"request_count={} request_paths={paths:?} server_task_finished={}",
paths.len(),
self.task.is_finished()
),
Err(error) => format!("request_paths_unavailable={error} server_task_finished={}", self.task.is_finished()),
}
}
}
impl Drop for TestServer {
fn drop(&mut self) {
self.task.abort();
@@ -678,7 +694,7 @@ async fn explicit_proxy_carries_registration_rotation_and_heartbeat_with_mtls()
.expect("configured heartbeat runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online { .. }
));
runtime.shutdown().await;
@@ -872,22 +888,33 @@ async fn wait_for_requests(server: &TestServer, count: usize) {
}
async fn wait_for_heartbeat_status(
server: &TestServer,
status: &mut watch::Receiver<HeartbeatStatus>,
predicate: impl Fn(&HeartbeatStatus) -> bool,
) -> HeartbeatStatus {
tokio::time::timeout(Duration::from_secs(3), async {
tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async {
loop {
let current = status.borrow_and_update().clone();
if predicate(&current) {
return current;
}
status.changed().await.unwrap_or_else(|error| {
panic!("heartbeat status channel: {error}; last status: {:?}", *status.borrow());
panic!(
"heartbeat status channel: {error}; last status: {:?}; {}",
*status.borrow(),
server.request_diagnostics()
);
});
}
})
.await
.unwrap_or_else(|error| panic!("heartbeat status: {error}; last status: {:?}", *status.borrow()))
.unwrap_or_else(|error| {
panic!(
"heartbeat status: {error}; last status: {:?}; {}",
*status.borrow(),
server.request_diagnostics()
)
})
}
fn rotation_response(pki: &TestPki, identity: &rustfs::connect::DeviceIdentity, serial: u8) -> (Value, Value) {
@@ -1301,7 +1328,25 @@ async fn heartbeat_runtime_respects_rotation_retry_after_without_blocking_heartb
.expect("configured runtime");
let mut status = runtime.status();
wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await;
// Start the progress bound after the durable rotation claim reaches the server.
// Heartbeats must still proceed before the five-second rotation Retry-After expires.
wait_for_requests(&server, 1).await;
assert!(
server.paths.lock().expect("paths lock")[0].ends_with(":rotateCredential"),
"rotation request must precede heartbeat progress"
);
tokio::time::timeout(
Duration::from_secs(3),
wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })),
)
.await
.unwrap_or_else(|error| {
panic!(
"heartbeat progress blocked after rotation: {error}; last status: {:?}; {}",
*status.borrow(),
server.request_diagnostics()
)
});
wait_for_requests(&server, 4).await;
runtime.shutdown().await;
@@ -1353,7 +1398,7 @@ async fn heartbeat_runtime_skips_only_valid_pending_reenrollment() {
let mut status = runtime.status();
assert!(matches!(
wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online { .. }
));
runtime.shutdown().await;
@@ -1393,7 +1438,8 @@ async fn heartbeat_runtime_skips_only_valid_pending_reenrollment() {
.expect("configured runtime");
let mut corrupted_status = corrupted_runtime.status();
assert_eq!(
wait_for_heartbeat_status(&mut corrupted_status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
wait_for_heartbeat_status(&server, &mut corrupted_status, |status| matches!(status, HeartbeatStatus::Failed { .. }))
.await,
HeartbeatStatus::Failed {
reason: HeartbeatError::StateConflict.to_string(),
}
@@ -1462,7 +1508,7 @@ async fn heartbeat_retries_with_the_new_credential_after_concurrent_rotation() {
.expect("concurrent rotation")
.expect("rotation due");
assert!(matches!(
wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online { .. }
));
runtime.shutdown().await;
@@ -1547,7 +1593,7 @@ async fn inventory_retries_with_the_new_credential_after_concurrent_rotation() {
.expect("concurrent rotation")
.expect("rotation due");
assert!(matches!(
tokio::time::timeout(Duration::from_secs(3), async {
tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async {
loop {
let current = status.borrow_and_update().clone();
if matches!(current, InventoryStatus::BackingOff { .. }) {
@@ -1642,20 +1688,37 @@ async fn inventory_first_recovers_a_saved_reenrollment_before_telemetry() {
.expect("start inventory runtime")
.expect("configured inventory runtime");
let mut inventory_status = inventory.status();
let startup_diagnostics = || {
format!(
"{}; pending_registration={} staged_identity={} completed_registration={}",
telemetry.request_diagnostics(),
temp.path().join("credential/registration.pending.json").exists(),
temp.path().join("identity/device.key.next").exists(),
temp.path().join("credential/registration.completed.json").exists()
)
};
assert!(matches!(
tokio::time::timeout(Duration::from_secs(3), async {
tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async {
loop {
let current = inventory_status.borrow_and_update().clone();
if matches!(current, InventoryStatus::Online { .. }) {
break current;
}
inventory_status.changed().await.unwrap_or_else(|error| {
panic!("inventory status channel: {error}; last status: {:?}", *inventory_status.borrow());
panic!(
"inventory status channel: {error}; last status: {:?}; {}",
*inventory_status.borrow(),
startup_diagnostics()
);
});
}
})
.await
.unwrap_or_else(|error| panic!("inventory online status: {error}; last status: {:?}", *inventory_status.borrow())),
.unwrap_or_else(|error| panic!(
"inventory online status: {error}; last status: {:?}; {}",
*inventory_status.borrow(),
startup_diagnostics()
)),
InventoryStatus::Online { .. }
));
@@ -1684,7 +1747,11 @@ async fn inventory_first_recovers_a_saved_reenrollment_before_telemetry() {
.expect("configured heartbeat runtime");
let mut heartbeat_status = heartbeat.status();
assert!(matches!(
wait_for_heartbeat_status(&mut heartbeat_status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
wait_for_heartbeat_status(&telemetry, &mut heartbeat_status, |status| matches!(
status,
HeartbeatStatus::Online { .. }
))
.await,
HeartbeatStatus::Online { .. }
));
heartbeat.shutdown().await;
@@ -1747,7 +1814,11 @@ async fn heartbeat_runtime_stops_when_rotation_reports_revocation() {
let mut status = runtime.status();
assert_eq!(
wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::AuthenticationStopped { .. })).await,
wait_for_heartbeat_status(&server, &mut status, |status| matches!(
status,
HeartbeatStatus::AuthenticationStopped { .. }
))
.await,
HeartbeatStatus::AuthenticationStopped {
status: 401,
reason: Some("DEVICE_REVOKED".to_owned()),