mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 04:16:38 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| aa727d1448 | |||
| f8ad979c7a | |||
| 941b6ab1d0 | |||
| cce3ecdbfc |
@@ -400,7 +400,7 @@ jobs:
|
|||||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||||
needs: [ quick-checks ]
|
needs: [ quick-checks ]
|
||||||
runs-on: sm-standard-4
|
runs-on: sm-standard-4
|
||||||
timeout-minutes: 90
|
timeout-minutes: 45
|
||||||
env:
|
env:
|
||||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||||
steps:
|
steps:
|
||||||
@@ -440,7 +440,7 @@ jobs:
|
|||||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||||
needs: [ quick-checks ]
|
needs: [ quick-checks ]
|
||||||
runs-on: sm-standard-4
|
runs-on: sm-standard-4
|
||||||
timeout-minutes: 90
|
timeout-minutes: 60
|
||||||
env:
|
env:
|
||||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||||
steps:
|
steps:
|
||||||
@@ -470,7 +470,7 @@ jobs:
|
|||||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||||
needs: [ quick-checks ]
|
needs: [ quick-checks ]
|
||||||
runs-on: sm-standard-4
|
runs-on: sm-standard-4
|
||||||
timeout-minutes: 90
|
timeout-minutes: 60
|
||||||
strategy:
|
strategy:
|
||||||
# On a PR, one failing protocol leg is enough to know the PR is not ready,
|
# On a PR, one failing protocol leg is enough to know the PR is not ready,
|
||||||
# so stop the sibling leg instead of paying another ~40 minutes for it.
|
# so stop the sibling leg instead of paying another ~40 minutes for it.
|
||||||
|
|||||||
@@ -57,13 +57,6 @@ pub const DEFAULT_MAX_IO_EVENTS_PER_TICK: usize = 1024;
|
|||||||
pub const DEFAULT_EVENT_INTERVAL: u32 = 61;
|
pub const DEFAULT_EVENT_INTERVAL: u32 = 61;
|
||||||
pub const DEFAULT_RNG_SEED: Option<u64> = None; // None means random
|
pub const DEFAULT_RNG_SEED: Option<u64> = None; // None means random
|
||||||
|
|
||||||
/// Dedicated blocking thread pool for fsync/fdatasync operations.
|
|
||||||
/// When > 1, fsync operations are isolated from the main blocking pool to
|
|
||||||
/// prevent device-bound fsync from starving read operations (pread/stat/open).
|
|
||||||
/// Default 0 means auto (no isolation, use main runtime).
|
|
||||||
pub const ENV_FSYNC_BLOCKING_THREADS: &str = "RUSTFS_RUNTIME_FSYNC_BLOCKING_THREADS";
|
|
||||||
pub const DEFAULT_FSYNC_BLOCKING_THREADS: usize = 0;
|
|
||||||
|
|
||||||
// Dial9 Tokio Telemetry Default values
|
// Dial9 Tokio Telemetry Default values
|
||||||
pub const DEFAULT_RUNTIME_DIAL9_ENABLED: bool = false; // Disabled by default
|
pub const DEFAULT_RUNTIME_DIAL9_ENABLED: bool = false; // Disabled by default
|
||||||
pub const DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR: &str = "/var/log/rustfs/telemetry";
|
pub const DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR: &str = "/var/log/rustfs/telemetry";
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ async fn test_bucket_default_sse_s3_put_object() -> Result<(), Box<dyn std::erro
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -159,7 +159,7 @@ async fn test_bucket_default_sse_kms_put_object() -> Result<(), Box<dyn std::err
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -278,7 +278,7 @@ async fn test_bucket_default_sse_kms_multipart_crc32() -> Result<(), Box<dyn std
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -475,7 +475,7 @@ async fn test_explicit_encryption_overrides_bucket_default() -> Result<(), Box<d
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -570,7 +570,7 @@ async fn test_sse_kms_without_key_id_populates_default() -> Result<(), Box<dyn s
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
|||||||
@@ -189,34 +189,69 @@ pub async fn wait_for_kms_ready(
|
|||||||
access_key: &str,
|
access_key: &str,
|
||||||
secret_key: &str,
|
secret_key: &str,
|
||||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||||
let total_deadline = Duration::from_secs(5);
|
wait_for_kms_ready_with_timeout(base_url, access_key, secret_key, Duration::from_secs(5)).await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn wait_for_kms_ready_with_timeout(
|
||||||
|
base_url: &str,
|
||||||
|
access_key: &str,
|
||||||
|
secret_key: &str,
|
||||||
|
total_deadline: Duration,
|
||||||
|
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||||
let start = tokio::time::Instant::now();
|
let start = tokio::time::Instant::now();
|
||||||
|
let deadline = start + total_deadline;
|
||||||
let mut backoff = Duration::from_millis(200);
|
let mut backoff = Duration::from_millis(200);
|
||||||
let max_backoff = Duration::from_secs(1);
|
let max_backoff = Duration::from_secs(1);
|
||||||
let mut first_attempt = true;
|
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
if !first_attempt {
|
match tokio::time::timeout_at(deadline, get_kms_status(base_url, access_key, secret_key)).await {
|
||||||
if start.elapsed() >= total_deadline {
|
Ok(Ok(status)) => {
|
||||||
return Err("KMS failed to become ready within 5 seconds".into());
|
|
||||||
}
|
|
||||||
sleep(backoff).await;
|
|
||||||
backoff = (backoff * 2).min(max_backoff);
|
|
||||||
}
|
|
||||||
first_attempt = false;
|
|
||||||
|
|
||||||
match get_kms_status(base_url, access_key, secret_key).await {
|
|
||||||
Ok(status) => {
|
|
||||||
info!("KMS is ready (status: {})", status);
|
info!("KMS is ready (status: {})", status);
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Ok(Err(e)) => {
|
||||||
if start.elapsed() >= total_deadline {
|
let elapsed_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX);
|
||||||
return Err(format!("KMS did not become ready within 5 s: last error: {e}").into());
|
warn!(error = %e, elapsed_ms, "KMS not ready yet, retrying…");
|
||||||
}
|
|
||||||
warn!(error = %e, elapsed_ms = start.elapsed().as_millis() as u64, "KMS not ready yet, retrying…");
|
|
||||||
}
|
}
|
||||||
|
Err(_) => return Err(format!("KMS failed to become ready within {} ms", total_deadline.as_millis()).into()),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let now = tokio::time::Instant::now();
|
||||||
|
if now >= deadline {
|
||||||
|
return Err(format!("KMS failed to become ready within {} ms", total_deadline.as_millis()).into());
|
||||||
|
}
|
||||||
|
sleep((now + backoff).min(deadline) - now).await;
|
||||||
|
backoff = (backoff * 2).min(max_backoff);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod readiness_tests {
|
||||||
|
use super::wait_for_kms_ready_with_timeout;
|
||||||
|
use std::time::Duration;
|
||||||
|
use tokio::io::AsyncReadExt;
|
||||||
|
use tokio::net::TcpListener;
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn kms_readiness_deadline_covers_a_stalled_status_request() {
|
||||||
|
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind readiness test server");
|
||||||
|
let address = listener.local_addr().expect("read readiness test server address");
|
||||||
|
let server = tokio::spawn(async move {
|
||||||
|
let (mut socket, _) = listener.accept().await.expect("accept readiness request");
|
||||||
|
let mut request = [0_u8; 1024];
|
||||||
|
let _ = socket.read(&mut request).await.expect("read readiness request");
|
||||||
|
std::future::pending::<()>().await;
|
||||||
|
});
|
||||||
|
|
||||||
|
let result = tokio::time::timeout(
|
||||||
|
Duration::from_secs(1),
|
||||||
|
wait_for_kms_ready_with_timeout(&format!("http://{address}"), "access-key", "secret-key", Duration::from_millis(50)),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("readiness helper must enforce its own deadline");
|
||||||
|
|
||||||
|
server.abort();
|
||||||
|
assert!(result.is_err(), "a stalled status request must not outlive the readiness deadline");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -61,7 +61,7 @@ async fn test_metadata_replace_self_copy_of_sse_object_stays_decryptable() {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("failed to start RustFS with local KMS");
|
.expect("failed to start RustFS with local KMS");
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
let client = kms_env.base_env.create_s3_client();
|
let client = kms_env.base_env.create_s3_client();
|
||||||
// Deliberately an UNVERSIONED bucket: that is the branch where the store layer can service
|
// Deliberately an UNVERSIONED bucket: that is the branch where the store layer can service
|
||||||
@@ -160,7 +160,7 @@ async fn test_metadata_replace_self_copy_dropping_sse_rewrites_plaintext() {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("failed to start RustFS with local KMS");
|
.expect("failed to start RustFS with local KMS");
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
let client = kms_env.base_env.create_s3_client();
|
let client = kms_env.base_env.create_s3_client();
|
||||||
// Unversioned, and deliberately WITHOUT a bucket default-encryption rule, so the copy below
|
// Unversioned, and deliberately WITHOUT a bucket default-encryption rule, so the copy below
|
||||||
@@ -256,7 +256,7 @@ async fn test_metadata_replace_self_copy_under_bucket_default_sse_stays_decrypta
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("failed to start RustFS with local KMS");
|
.expect("failed to start RustFS with local KMS");
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
let client = kms_env.base_env.create_s3_client();
|
let client = kms_env.base_env.create_s3_client();
|
||||||
let bucket = "copy-object-self-copy-bucket-default-sse-test";
|
let bucket = "copy-object-self-copy-bucket-default-sse-test";
|
||||||
|
|||||||
@@ -56,7 +56,7 @@ async fn test_self_copy_of_historical_sse_s3_version_is_readable() {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("failed to start RustFS with local KMS");
|
.expect("failed to start RustFS with local KMS");
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
let client = kms_env.base_env.create_s3_client();
|
let client = kms_env.base_env.create_s3_client();
|
||||||
let bucket = "copy-object-version-restore-sse-test";
|
let bucket = "copy-object-version-restore-sse-test";
|
||||||
|
|||||||
@@ -87,7 +87,7 @@ async fn test_head_reports_managed_metadata_for_sse_s3() -> Result<(), Box<dyn s
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -147,7 +147,7 @@ async fn test_head_reports_managed_metadata_for_sse_kms_and_copy() -> Result<(),
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -250,7 +250,7 @@ async fn test_multipart_upload_writes_encrypted_data() -> Result<(), Box<dyn std
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
|||||||
@@ -24,7 +24,6 @@ use super::common::{
|
|||||||
test_sse_kms_encryption, test_sse_s3_encryption,
|
test_sse_kms_encryption, test_sse_s3_encryption,
|
||||||
};
|
};
|
||||||
use crate::common::{TEST_BUCKET, init_logging};
|
use crate::common::{TEST_BUCKET, init_logging};
|
||||||
use tokio::time::{Duration, sleep};
|
|
||||||
use tracing::info;
|
use tracing::info;
|
||||||
|
|
||||||
/// Comprehensive test: Full KMS workflow with all encryption types
|
/// Comprehensive test: Full KMS workflow with all encryption types
|
||||||
@@ -35,7 +34,7 @@ async fn test_comprehensive_kms_full_workflow() -> Result<(), Box<dyn std::error
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -103,7 +102,7 @@ async fn test_comprehensive_stress_test() -> Result<(), Box<dyn std::error::Erro
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -137,7 +136,7 @@ async fn test_comprehensive_key_isolation() -> Result<(), Box<dyn std::error::Er
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -208,7 +207,7 @@ async fn test_comprehensive_concurrent_operations() -> Result<(), Box<dyn std::e
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -253,7 +252,7 @@ async fn test_comprehensive_performance_benchmark() -> Result<(), Box<dyn std::e
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
|||||||
@@ -44,7 +44,7 @@ async fn test_kms_zero_byte_file_encryption() -> Result<(), Box<dyn std::error::
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -117,7 +117,7 @@ async fn test_kms_single_byte_file_encryption() -> Result<(), Box<dyn std::error
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -209,7 +209,7 @@ async fn test_kms_multipart_boundary_conditions() -> Result<(), Box<dyn std::err
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -284,7 +284,7 @@ async fn test_kms_invalid_key_scenarios() -> Result<(), Box<dyn std::error::Erro
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -371,7 +371,7 @@ async fn test_kms_concurrent_encryption() -> Result<(), Box<dyn std::error::Erro
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = Arc::new(kms_env.base_env.create_s3_client());
|
let s3_client = Arc::new(kms_env.base_env.create_s3_client());
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -478,7 +478,7 @@ async fn test_kms_key_validation_security() -> Result<(), Box<dyn std::error::Er
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ async fn test_kms_key_directory_unavailable() -> Result<(), Box<dyn std::error::
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -127,7 +127,7 @@ async fn test_kms_corrupted_key_files() -> Result<(), Box<dyn std::error::Error
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -218,7 +218,7 @@ async fn test_kms_multipart_upload_interruption() -> Result<(), Box<dyn std::err
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -401,7 +401,7 @@ async fn test_kms_resource_constraints() -> Result<(), Box<dyn std::error::Error
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
|||||||
@@ -46,7 +46,7 @@ async fn test_local_kms_end_to_end() -> Result<(), Box<dyn std::error::Error + S
|
|||||||
.expect("Failed to start RustFS with Local KMS");
|
.expect("Failed to start RustFS with Local KMS");
|
||||||
|
|
||||||
// Wait a moment for RustFS to fully start up and initialize KMS
|
// Wait a moment for RustFS to fully start up and initialize KMS
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
||||||
|
|
||||||
@@ -127,7 +127,7 @@ async fn test_local_kms_key_isolation() {
|
|||||||
.expect("Failed to start RustFS with Local KMS");
|
.expect("Failed to start RustFS with Local KMS");
|
||||||
|
|
||||||
// Wait a moment for RustFS to fully start up and initialize KMS
|
// Wait a moment for RustFS to fully start up and initialize KMS
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
||||||
|
|
||||||
@@ -227,7 +227,7 @@ async fn test_local_kms_large_file() {
|
|||||||
.expect("Failed to start RustFS with Local KMS");
|
.expect("Failed to start RustFS with Local KMS");
|
||||||
|
|
||||||
// Wait a moment for RustFS to fully start up and initialize KMS
|
// Wait a moment for RustFS to fully start up and initialize KMS
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
||||||
|
|
||||||
@@ -309,7 +309,7 @@ async fn test_local_kms_multipart_upload() {
|
|||||||
.expect("Failed to start RustFS with Local KMS");
|
.expect("Failed to start RustFS with Local KMS");
|
||||||
|
|
||||||
// Wait for KMS initialization
|
// Wait for KMS initialization
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await.expect("KMS ready");
|
||||||
|
|
||||||
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
info!("RustFS started with KMS auto-configuration, default_key_id: {}", default_key_id);
|
||||||
|
|
||||||
|
|||||||
@@ -19,7 +19,6 @@
|
|||||||
//! multipart upload behaviour.
|
//! multipart upload behaviour.
|
||||||
|
|
||||||
use crate::common::{TEST_BUCKET, init_logging};
|
use crate::common::{TEST_BUCKET, init_logging};
|
||||||
use tokio::time::{Duration, sleep};
|
|
||||||
use tracing::{error, info};
|
use tracing::{error, info};
|
||||||
|
|
||||||
use super::common::{
|
use super::common::{
|
||||||
@@ -45,8 +44,8 @@ impl VaultKmsTestContext {
|
|||||||
|
|
||||||
start_kms(&env.base_env.url, &env.base_env.access_key, &env.base_env.secret_key).await?;
|
start_kms(&env.base_env.url, &env.base_env.access_key, &env.base_env.secret_key).await?;
|
||||||
|
|
||||||
// Allow Vault to finish initialising token auth and transit engine.
|
// Wait for KMS to finish initialising.
|
||||||
sleep(Duration::from_secs(2)).await;
|
super::common::wait_for_kms_ready(&env.base_env.url, &env.base_env.access_key, &env.base_env.secret_key).await?;
|
||||||
|
|
||||||
Ok(Self { env })
|
Ok(Self { env })
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -33,7 +33,7 @@ async fn test_step1_basic_single_file_encryption() -> Result<(), Box<dyn std::er
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -89,7 +89,7 @@ async fn test_step2_basic_multipart_upload_without_encryption() -> Result<(), Bo
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -187,7 +187,7 @@ async fn test_step3_multipart_upload_with_sse_s3() -> Result<(), Box<dyn std::er
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -310,7 +310,7 @@ async fn test_step4_large_multipart_upload_with_encryption() -> Result<(), Box<d
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
@@ -435,7 +435,7 @@ async fn test_step5_all_encryption_types_multipart() -> Result<(), Box<dyn std::
|
|||||||
|
|
||||||
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
let mut kms_env = LocalKMSTestEnvironment::new().await?;
|
||||||
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
let _default_key_id = kms_env.start_rustfs_for_local_kms().await?;
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(3)).await;
|
kms_env.wait_for_kms_ready().await?;
|
||||||
|
|
||||||
let s3_client = kms_env.base_env.create_s3_client();
|
let s3_client = kms_env.base_env.create_s3_client();
|
||||||
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
kms_env.base_env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
|||||||
@@ -315,7 +315,7 @@ pub async fn fsync_dir(dir: impl AsRef<Path>) -> io::Result<()> {
|
|||||||
#[cfg(unix)]
|
#[cfg(unix)]
|
||||||
{
|
{
|
||||||
let dir = dir.as_ref().to_path_buf();
|
let dir = dir.as_ref().to_path_buf();
|
||||||
fsync_spawn_blocking(move || fsync_dir_std(dir)).await?
|
tokio::task::spawn_blocking(move || fsync_dir_std(dir)).await?
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(not(unix))]
|
#[cfg(not(unix))]
|
||||||
@@ -683,7 +683,7 @@ async fn fsync_open_dst_dir_group(group: &DstDirFsyncGroup) -> io::Result<()> {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
let dir = group.dir.clone();
|
let dir = group.dir.clone();
|
||||||
let dir_file = group.dir_file.clone();
|
let dir_file = group.dir_file.clone();
|
||||||
fsync_spawn_blocking(move || {
|
tokio::task::spawn_blocking(move || {
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
{
|
{
|
||||||
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
|
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
|
||||||
@@ -1080,44 +1080,6 @@ const TEST_GLOBAL_FILE_SYNCS: usize = 64;
|
|||||||
|
|
||||||
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
||||||
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||||
|
|
||||||
/// Dedicated tokio runtime for fsync/fdatasync blocking operations. When
|
|
||||||
/// configured with >1 threads, isolates device-bound fsync from the main
|
|
||||||
/// blocking pool so reads (pread/stat/open) are not starved. `None` means
|
|
||||||
/// fall back to the main runtime (zero behavior change).
|
|
||||||
static FSYNC_RUNTIME: LazyLock<Option<tokio::runtime::Runtime>> = LazyLock::new(|| {
|
|
||||||
let threads =
|
|
||||||
rustfs_utils::get_env_usize(rustfs_config::ENV_FSYNC_BLOCKING_THREADS, rustfs_config::DEFAULT_FSYNC_BLOCKING_THREADS);
|
|
||||||
if threads <= 1 {
|
|
||||||
return None;
|
|
||||||
}
|
|
||||||
let mut builder = tokio::runtime::Builder::new_multi_thread();
|
|
||||||
builder
|
|
||||||
.worker_threads(num_cpus::get().min(8))
|
|
||||||
.max_blocking_threads(threads)
|
|
||||||
.thread_name("rustfs-fsync")
|
|
||||||
.thread_stack_size(512 * 1024)
|
|
||||||
.enable_all();
|
|
||||||
match builder.build() {
|
|
||||||
Ok(rt) => {
|
|
||||||
tracing::info!(threads, "fsync dedicated blocking pool enabled");
|
|
||||||
Some(rt)
|
|
||||||
}
|
|
||||||
Err(err) => {
|
|
||||||
tracing::warn!(%err, "failed to build fsync runtime, falling back to main pool");
|
|
||||||
None
|
|
||||||
}
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
/// Spawn a blocking task on the fsync-dedicated runtime if configured,
|
|
||||||
/// otherwise fall back to the main tokio blocking pool.
|
|
||||||
fn fsync_spawn_blocking<T: Send + 'static>(f: impl FnOnce() -> T + Send + 'static) -> tokio::task::JoinHandle<T> {
|
|
||||||
match FSYNC_RUNTIME.as_ref() {
|
|
||||||
Some(rt) => rt.spawn_blocking(f),
|
|
||||||
None => tokio::task::spawn_blocking(f),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
||||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||||
type NamespaceMutationLock = AsyncMutex<()>;
|
type NamespaceMutationLock = AsyncMutex<()>;
|
||||||
@@ -1255,7 +1217,7 @@ where
|
|||||||
F: FnOnce() -> io::Result<T> + Send + 'static,
|
F: FnOnce() -> io::Result<T> + Send + 'static,
|
||||||
{
|
{
|
||||||
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
|
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
|
||||||
let result = fsync_spawn_blocking(move || {
|
let result = tokio::task::spawn_blocking(move || {
|
||||||
let _disk_permit = disk_permit;
|
let _disk_permit = disk_permit;
|
||||||
work()
|
work()
|
||||||
})
|
})
|
||||||
@@ -2184,7 +2146,7 @@ async fn run_blocking_namespace_file_sync_operation_with_global<T: Send + 'stati
|
|||||||
wait_started,
|
wait_started,
|
||||||
);
|
);
|
||||||
let disk_permit = admission.disk_permit.clone();
|
let disk_permit = admission.disk_permit.clone();
|
||||||
let result = fsync_spawn_blocking(move || {
|
let result = tokio::task::spawn_blocking(move || {
|
||||||
let _lease = lease;
|
let _lease = lease;
|
||||||
let _disk_permit = disk_permit;
|
let _disk_permit = disk_permit;
|
||||||
operation()
|
operation()
|
||||||
|
|||||||
@@ -206,13 +206,6 @@ def check_runner_selection(root: Path) -> list[str]:
|
|||||||
return errors
|
return errors
|
||||||
|
|
||||||
|
|
||||||
def check_s3_tests_runner(root: Path) -> list[str]:
|
|
||||||
runner = (root / "scripts/s3-tests/run.sh").read_text()
|
|
||||||
if "--showlocals" in runner:
|
|
||||||
return ["scripts/s3-tests/run.sh: pytest failure diagnostics must not dump local values"]
|
|
||||||
return []
|
|
||||||
|
|
||||||
|
|
||||||
def profile_selection(root: Path, profile: str) -> str:
|
def profile_selection(root: Path, profile: str) -> str:
|
||||||
if not re.fullmatch(r"e2e-[a-z0-9-]+", profile):
|
if not re.fullmatch(r"e2e-[a-z0-9-]+", profile):
|
||||||
raise ValueError(f"invalid e2e profile name: {profile}")
|
raise ValueError(f"invalid e2e profile name: {profile}")
|
||||||
@@ -279,7 +272,6 @@ def validate(root: Path) -> list[str]:
|
|||||||
errors.extend(check_e2e_modules(root))
|
errors.extend(check_e2e_modules(root))
|
||||||
errors.extend(check_fuzz_targets(root))
|
errors.extend(check_fuzz_targets(root))
|
||||||
errors.extend(check_runner_selection(root))
|
errors.extend(check_runner_selection(root))
|
||||||
errors.extend(check_s3_tests_runner(root))
|
|
||||||
errors.extend(check_profile_definitions(root))
|
errors.extend(check_profile_definitions(root))
|
||||||
return errors
|
return errors
|
||||||
|
|
||||||
@@ -349,23 +341,6 @@ class SelfTests(unittest.TestCase):
|
|||||||
)
|
)
|
||||||
self.assertEqual(len(check_fuzz_targets(root)), 1)
|
self.assertEqual(len(check_fuzz_targets(root)), 1)
|
||||||
|
|
||||||
def test_s3_runner_rejects_unbounded_failure_locals(self) -> None:
|
|
||||||
with tempfile.TemporaryDirectory() as tmp:
|
|
||||||
root = Path(tmp)
|
|
||||||
runner = root / "scripts/s3-tests/run.sh"
|
|
||||||
runner.parent.mkdir(parents=True)
|
|
||||||
runner.write_text("tox -- -vv -ra --tb=long\n")
|
|
||||||
self.assertEqual(check_s3_tests_runner(root), [])
|
|
||||||
runner.write_text("tox -- -vv -ra --showlocals --tb=long\n")
|
|
||||||
self.assertEqual(len(check_s3_tests_runner(root)), 1)
|
|
||||||
with (
|
|
||||||
mock.patch(__name__ + ".check_e2e_modules", return_value=[]),
|
|
||||||
mock.patch(__name__ + ".check_fuzz_targets", return_value=[]),
|
|
||||||
mock.patch(__name__ + ".check_runner_selection", return_value=[]),
|
|
||||||
mock.patch(__name__ + ".check_profile_definitions", return_value=[]),
|
|
||||||
):
|
|
||||||
self.assertEqual(len(validate(root)), 1)
|
|
||||||
|
|
||||||
def test_profile_listing_enforces_selection(self) -> None:
|
def test_profile_listing_enforces_selection(self) -> None:
|
||||||
with tempfile.TemporaryDirectory() as tmp:
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
root = Path(tmp)
|
root = Path(tmp)
|
||||||
@@ -436,7 +411,7 @@ def main() -> int:
|
|||||||
for error in errors:
|
for error in errors:
|
||||||
print(f"ERROR: {error}", file=sys.stderr)
|
print(f"ERROR: {error}", file=sys.stderr)
|
||||||
return 1
|
return 1
|
||||||
print("OK: e2e modules, runner selection, fuzz matrices, profiles, and bounded diagnostics are wired")
|
print("OK: e2e modules, runner selection, fuzz matrices, and profile guards are wired")
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1028,11 +1028,10 @@ else
|
|||||||
fi
|
fi
|
||||||
|
|
||||||
# Run tests from s3tests/functional
|
# Run tests from s3tests/functional
|
||||||
# Failure locals can contain multi-MiB request bodies; keep tracebacks without expanding local values.
|
|
||||||
set +e
|
set +e
|
||||||
S3TEST_CONF="${CONF_OUTPUT_PATH}" \
|
S3TEST_CONF="${CONF_OUTPUT_PATH}" \
|
||||||
tox -- \
|
tox -- \
|
||||||
-vv -ra --tb=long \
|
-vv -ra --showlocals --tb=long \
|
||||||
--maxfail="${MAXFAIL}" \
|
--maxfail="${MAXFAIL}" \
|
||||||
--timeout="${TEST_TIMEOUT}" \
|
--timeout="${TEST_TIMEOUT}" \
|
||||||
--junitxml="${ARTIFACTS_DIR}/junit.xml" \
|
--junitxml="${ARTIFACTS_DIR}/junit.xml" \
|
||||||
|
|||||||
Reference in New Issue
Block a user