fix(storage): weight automatic multipart admission by part size (#8118)

* fix(storage): weight automatic multipart admission by part size

* fix(ci): restore filesystem runner capabilities and typos dependency

* fix(ci): make release guard portable and spell out part variables

* ci: restore sm-standard-2 runners for io_uring and distributed e2e jobs

---------

Co-authored-by: hector <42570491+majinghe@users.noreply.github.com>
Co-authored-by: Hauser <housemecn@gmail.com>
This commit is contained in:
Chris
2026-09-26 22:35:39 +08:00
committed by GitHub
parent 212b20f890
commit 8741bb77ea
10 changed files with 633 additions and 58 deletions
+2 -2
View File
@@ -96,11 +96,11 @@ pub struct WorkloadAdmissionSnapshot {
pub class: WorkloadClass,
/// Current admission state.
pub state: AdmissionState,
/// Active work count when the owner exposes it.
/// Active work count or allocated/reserved permit units, as defined by the owner.
pub active: Option<usize>,
/// Queued work count when the owner exposes it.
pub queued: Option<usize>,
/// Admission limit when the owner exposes it.
/// Admission limit in the same units as `active`, when the owner exposes it.
pub limit: Option<usize>,
/// Optional state reason for disabled, throttled, saturated, or unknown states.
pub reason: Option<String>,
+11 -9
View File
@@ -169,15 +169,17 @@ Scanner cycle budget controls:
## Foreground write admission environment variables
Large direct `PutObject` requests and multipart `UploadPart` requests share one
per-process permit pool that bounds how many bodies are ingested and written
concurrently. Small direct PUTs stay on the legacy path.
per-process permit pool that bounds concurrent body ingest and storage writes.
Small direct PUTs stay on the legacy path.
- `RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE`
- enables the default-on pool; `false` keeps only the soft request counter.
- default is `true`.
- `RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT`
- permits in the pool; `0` derives half of `RUSTFS_OBJECT_MAX_CONCURRENT_DISK_READS`, clamped to `32`.
- default is `0` (32 permits at stock settings).
- a positive value is an exact request-count limit shared by all gated writes.
- default is `0`: derive half of `RUSTFS_OBJECT_MAX_CONCURRENT_DISK_READS`, clamped to `32`, as large-write slots. Each slot has eight units. A gated direct PUT or unknown-size multipart part uses eight units; known-size parts use one unit per 8 MiB, rounded up, with a minimum of one and maximum of eight.
- stock settings therefore share 256 units across at most 32 large/unknown writes or 256 parts of up to 8 MiB. Mixed sizes consume the same budget. Admitted known parts of up to 64 MiB total at most 2 GiB of declared payload; larger streaming writes and queued bodies are separate. This is not a bound on total process or kernel memory.
- admission snapshots report allocated or reserved units in `active` and the unit budget in `limit` under automatic sizing. Explicit and strict limits still report request-count permits.
- `RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES`
- smallest direct `PutObject` that takes a permit; unknown-size requests always do.
- default is `33554432` (32 MiB).
@@ -193,8 +195,8 @@ concurrently. Small direct PUTs stay on the legacy path.
- RustFS does not read the request body while a part is queued, so the client's socket write stalls for the whole wait and whatever timeout the client or an intermediary has configured competes with this value. Keep it with margin below the shortest such timeout in use (botocore applies its 60 s `connect_timeout` to the body write; the AWS SDK for Java v2 has a 30 s socket write timeout; reverse proxies add their own body timeouts); a wait that outlives the client timeout surfaces as a dropped connection instead of `SlowDown`.
- `RUSTFS_PUT_MULTIPART_FOREGROUND_ADMISSION_MAX_PENDING`
- maximum `UploadPart` requests waiting for a permit at once; parts beyond it return `SlowDown` without waiting.
- default is `0`, which derives 16 times the permit limit (512 at stock settings).
- each queued HTTP/1 part holds whatever unread body the client already pushed into the connection's kernel receive buffer (an HTTP/2 part holds up to its flow-control window in process memory), so this depth also bounds that memory. RustFS leaves the receive buffer to kernel autotuning (see `RUSTFS_HTTP_SOCKET_RECV_BUFFER_BYTES` below), which keeps an unread connection at the kernel's initial size (128 KiB on current Linux).
- default is `0`, which derives 16 times the large-write slot count, or the explicit request-count limit (512 at stock settings). Automatic subdivision does not enlarge this queue.
- each queued HTTP/1 part holds whatever unread body the client already pushed into the connection's kernel receive buffer (an HTTP/2 part holds up to its flow-control window in process memory). Autotuned receive buffers can remain large on reused connections; queue depth does not imply a fixed per-connection memory cost.
- `RUSTFS_PUT_FOREGROUND_ADMISSION_ENABLE`, `RUSTFS_PUT_FOREGROUND_ADMISSION_LIMIT`, `RUSTFS_PUT_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS`
- experimental strict gate that applies to every foreground write regardless of size and replaces the pool above when enabled.
- default is disabled; enabling it with limit `0` disables foreground write admission entirely.
@@ -203,9 +205,9 @@ concurrently. Small direct PUTs stay on the legacy path.
- `RUSTFS_HTTP_SOCKET_RECV_BUFFER_BYTES`
- fixed `SO_RCVBUF` for the API listener, inherited by every accepted socket; `0` leaves the receive buffer to kernel autotuning.
- default is `0`. Earlier releases hard-coded 4 MiB, which Linux doubles to 8 MiB and which disables autotuning, so every connection whose body was not being read yet (a multipart part queued for a foreground write permit) could accumulate up to 8 MiB of unread body in kernel memory; at SDK-default multipart concurrency that was enough to push a node into TCP memory pressure.
- with autotuning the per-connection receive ceiling is the kernel's (`net.ipv4.tcp_rmem` max, 6 MiB on stock Linux) instead of the former fixed 8 MiB, so a single very high-bandwidth-delay connection may see a somewhat lower ceiling; raise `net.ipv4.tcp_rmem` first, and set this variable only on kernels without receive-buffer autotuning (illumos/Solaris) or where the sysctl cannot be changed.
- the send buffer stays fixed at 4 MiB because the stock Linux send autotuning ceiling (`net.ipv4.tcp_wmem` max, 4 MiB) is lower than a GB-level response stream needs.
- default is `0`. A fixed buffer disables receive autotuning. Linux doubles the requested value for socket-memory accounting, subject to kernel limits; that capacity is not a measurement of queued payload or allocated memory.
- with autotuning the receive ceiling is controlled by the kernel (`net.ipv4.tcp_rmem` on Linux); its defaults vary with kernel version and memory. A connection can retain a buffer enlarged by previous requests while its next request is queued. Tune the kernel ceiling first, and set this variable only on kernels without receive-buffer autotuning (illumos/Solaris) or where the sysctl cannot be changed.
- the send buffer remains fixed at 4 MiB; this setting only controls the receive side.
## Remote tier timeout environment variables
+7 -4
View File
@@ -364,11 +364,14 @@ const _: () = assert!(!DEFAULT_PUT_FOREGROUND_ADMISSION_ENABLE);
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE";
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: bool = true;
/// Maximum automatic foreground write requests admitted concurrently per process.
/// Explicit maximum foreground write requests admitted concurrently per process.
///
/// `0` derives a conservative default from the local disk-read scheduler cap,
/// currently clamped to protect the commit path without making ordinary high
/// throughput uploads single-file.
/// `0` derives up to 32 large-write slots from the local disk-read scheduler.
/// Each automatic slot has eight units: a gated direct PUT or unknown-size
/// part uses all eight, while a known-size part uses one unit per 8 MiB,
/// rounded up and capped at eight. Thus stock settings share one budget across
/// at most 32 large/unknown writes or 256 parts of up to 8 MiB. A positive
/// override retains request-count semantics for every gated write.
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT";
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT: usize = 0;
+4 -7
View File
@@ -164,13 +164,10 @@ pub const DEFAULT_HTTP1_MAX_BUF_SIZE: usize = 64 * 1024; // 64 KB
/// autotuning.
///
/// A fixed `SO_RCVBUF` is inherited by every accepted socket and disables
/// receive-buffer autotuning, so a connection whose request body is not being
/// read yet (a multipart part queued for a foreground write permit) lets up to
/// the fixed size of unread body accumulate in kernel memory — Linux doubles
/// the requested value, so the former hard-coded 4 MiB held up to 8 MiB per
/// queued connection (issue #7385). Autotuning keeps an unread connection at
/// the kernel's initial size and grows only connections that are being
/// drained. Set this only on kernels without receive-buffer autotuning
/// receive-buffer autotuning. Queued multipart bodies can accumulate in kernel
/// memory with either policy: a reused autotuned connection may already have
/// a large buffer. Socket capacity is not the same as queued payload or actual
/// memory allocation. Set this only on kernels without receive-buffer autotuning
/// (illumos/Solaris) or on very high-bandwidth-delay links where the kernel's
/// autotuning ceiling (`net.ipv4.tcp_rmem` on Linux) is too low and cannot be
/// raised.
@@ -12,7 +12,7 @@
| Class | Provider (`impl WorkloadAdmissionSnapshotProvider`) | `active` / `queued` / `limit` source | Reports `Unknown` when |
|---|---|---|---|
| `ForegroundRead` | `ConcurrencyManager` in `rustfs/src/storage/concurrency/manager.rs` (source of truth); re-exposed unchanged by the RustFS runtime provider | disk-read permits in use / `None` (the semaphore exposes no waiter count) / configured max concurrent disk reads | the storage registry has no entry |
| `ForegroundWrite` | `ConcurrencyManager` in `rustfs/src/storage/concurrency/manager.rs` (source of truth); re-exposed unchanged by the RustFS runtime provider | foreground-write permits in use or legacy active-write counter / multipart parts waiting in the bounded admission queue (`None` for the strict and legacy policies) / configured or derived write-admission limit | the storage registry has no entry |
| `ForegroundWrite` | `ConcurrencyManager` in `rustfs/src/storage/concurrency/manager.rs` (source of truth); re-exposed unchanged by the RustFS runtime provider | allocated/reserved permit units under automatic sizing, request-count permits for explicit/strict limits, or the legacy active-write counter / multipart requests waiting in the bounded admission queue (`None` for the strict and legacy policies) / matching unit budget or request-count limit | the storage registry has no entry |
| `Metadata` | `RustFsWorkloadAdmissionSnapshotProvider` in `rustfs/src/workload_admission.rs` | `Open` once the bucket metadata runtime handle exists; no counts | bucket metadata runtime not initialized |
| `Scanner` | same | scanner active work-unit counter / none / configured set-scan limit when nonzero | scanner runtime not initialized |
| `Repair` | same | heal active tasks / heal queue length / `None` (limits live behind the async heal manager state) | heal manager not initialized |
@@ -50,7 +50,7 @@ Drive each cell with the stage histograms emitted by `crates/io-metrics/src/lib.
| Metric | Labels | Use |
| --- | --- | --- |
| `rustfs_s3_put_object_stage_duration_ms` | `stage` | P50/P95/P99 per PUT stage (`app_*`, `ingress_prepare`, `set_disk_*`) |
| `rustfs_s3_put_object_stage_duration_ms` | `stage` | P50/P95/P99 per PUT stage (`app_*`, `ingress_prepare`, `set_disk_*`, `multipart_*`) |
| `rustfs_io_get_object_stage_duration_seconds` | `path`, `stage` | Per GET stage, split by read path (`legacy_duplex`, `codec_streaming`, ...) |
| `rustfs_ec_encode_inflight_bytes_current` | — | EC encode memory pressure; pair with node RSS and CPU |
+9 -7
View File
@@ -1182,10 +1182,10 @@ impl DefaultMultipartUsecase {
{
return Err(S3Error::new(S3ErrorCode::EntityTooLarge));
}
let upload_part_admission = match self
.concurrency_manager()
.admit_multipart_part(size)
.await
let admission_wait_started = rustfs_io_metrics::put_stage_timer();
let admission = self.concurrency_manager().admit_multipart_part(size).await;
rustfs_io_metrics::record_put_object_stage_duration_from("multipart_admission_wait", admission_wait_started);
let upload_part_admission = match admission
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "foreground write admission closed"))?
{
ForegroundWriteAdmission::Disabled => None,
@@ -1385,11 +1385,13 @@ impl DefaultMultipartUsecase {
}
let _upload_part_admission = upload_part_admission;
let info = store
let store_write_started = rustfs_io_metrics::put_stage_timer();
let result = store
.put_object_part(&bucket, &key, &upload_id, part_id, &mut reader, &opts)
.await
.map_err(ApiError::from)?;
.await;
rustfs_io_metrics::record_put_object_stage_duration_from("multipart_store_write", store_write_started);
drop(_upload_part_admission);
let info = result.map_err(ApiError::from)?;
let mut checksum_crc32 = input.checksum_crc32;
let mut checksum_crc32c = input.checksum_crc32c;
+6 -12
View File
@@ -1170,18 +1170,12 @@ pub async fn start_http_server(
// 4. Socket buffers. The receive buffer is left to kernel autotuning
// unless RUSTFS_HTTP_SOCKET_RECV_BUFFER_BYTES is set: a fixed SO_RCVBUF
// is inherited by every accepted socket and disables autotuning, so a
// request whose body is not being read yet (a multipart part queued
// for a foreground write permit) lets up to the fixed size of unread
// body accumulate in kernel memory — the former hard-coded 4 MiB held
// up to 8 MiB per queued connection on Linux, which doubles the
// requested size. Autotuning keeps an unread connection at the
// kernel's initial size and grows only connections that are actually
// being drained (issue #7385). The send buffer stays fixed at 4 MiB
// because the stock Linux send autotuning ceiling (`tcp_wmem` max,
// 4 MiB) is below what a GB-level response stream needs, whereas the
// receive ceiling (`tcp_rmem` max, 6 MiB) already exceeds the old
// fixed request.
// is inherited by every accepted socket and disables autotuning.
// A reused autotuned connection can retain a buffer enlarged by a
// previous request, so queued multipart bodies still consume kernel
// memory. Neither policy makes socket capacity an allocation or
// guarantees a fixed per-waiter memory footprint. Keep the existing
// send-buffer tuning independent of receive autotuning.
// Some constrained local environments reject these socket options with
// EPERM/ENOPROTOOPT-style failures; log and continue in that case.
if recv_buffer_bytes > 0
+321 -15
View File
@@ -36,6 +36,12 @@ use tokio::sync::{OwnedSemaphorePermit, Semaphore};
use tracing::debug;
const DERIVED_LARGE_PUT_ADMISSION_LIMIT_MAX: usize = 32;
// The automatically sized pool keeps the large-write limit, but splits each
// slot so short multipart bodies do not occupy the same budget as long writes.
// At stock settings this admits at most 32 large/unknown writes or 256 parts of
// up to 8 MiB, sharing one FIFO budget. Explicit request-count limits stay exact.
const AUTO_FOREGROUND_WRITE_PERMITS_PER_SLOT: u16 = 8;
const MULTIPART_ADMISSION_UNIT_BYTES: u64 = 8 * 1024 * 1024;
// A queued multipart part holds a connection but no user-space body buffer on
// HTTP/1 (only whatever unread body the client already pushed into the kernel
// receive buffer); on HTTP/2 it holds up to the per-stream flow-control window
@@ -206,8 +212,9 @@ impl ForegroundWriteAdmissionGate {
&self,
wait_timeout: Duration,
max_pending: usize,
permits: u32,
) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
match self.semaphore.clone().try_acquire_owned() {
match self.semaphore.clone().try_acquire_many_owned(permits) {
Ok(permit) => return Ok(ForegroundWriteAdmission::Admitted(permit)),
Err(tokio::sync::TryAcquireError::Closed) => return Ok(ForegroundWriteAdmission::Rejected),
Err(tokio::sync::TryAcquireError::NoPermits) => {}
@@ -218,22 +225,22 @@ impl ForegroundWriteAdmissionGate {
let Some(_slot) = PendingSlot::reserve(&self.pending, max_pending) else {
return Ok(ForegroundWriteAdmission::Rejected);
};
match tokio::time::timeout(wait_timeout, self.semaphore.clone().acquire_owned()).await {
match tokio::time::timeout(wait_timeout, self.semaphore.clone().acquire_many_owned(permits)).await {
Ok(permit) => Ok(ForegroundWriteAdmission::Admitted(permit?)),
Err(_) => Ok(ForegroundWriteAdmission::Rejected),
}
}
async fn admit(&self) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
async fn admit(&self, permits: u32) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
if self.wait_timeout.is_zero() {
return Ok(match self.semaphore.clone().try_acquire_owned() {
return Ok(match self.semaphore.clone().try_acquire_many_owned(permits) {
Ok(permit) => ForegroundWriteAdmission::Admitted(permit),
Err(tokio::sync::TryAcquireError::NoPermits) => ForegroundWriteAdmission::Rejected,
Err(tokio::sync::TryAcquireError::Closed) => ForegroundWriteAdmission::Rejected,
});
}
match tokio::time::timeout(self.wait_timeout, self.semaphore.clone().acquire_owned()).await {
match tokio::time::timeout(self.wait_timeout, self.semaphore.clone().acquire_many_owned(permits)).await {
Ok(permit) => Ok(ForegroundWriteAdmission::Admitted(permit?)),
Err(_) => Ok(ForegroundWriteAdmission::Rejected),
}
@@ -252,6 +259,7 @@ enum ForegroundWriteAdmissionPolicy {
/// Default foreground write admission gate for pressure-heavy writes.
Large {
gate: ForegroundWriteAdmissionGate,
large_request_permits: u32,
put_object_min_size_bytes: usize,
multipart_part_min_size_bytes: usize,
multipart_wait_timeout: Duration,
@@ -295,13 +303,18 @@ impl ForegroundWriteAdmissionPolicy {
return Self::LegacyCounterOnly;
}
let large_limit = derive_large_put_admission_limit(
rustfs_utils::get_env_usize(
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT,
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT,
),
max_disk_reads,
let configured_limit = rustfs_utils::get_env_usize(
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT,
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT,
);
let large_limit = derive_large_put_admission_limit(configured_limit, max_disk_reads);
let large_request_permits = if configured_limit == 0 {
AUTO_FOREGROUND_WRITE_PERMITS_PER_SLOT
} else {
1
};
// The derived limit is clamped to 32; an explicit limit uses one permit.
let permit_limit = large_limit * usize::from(large_request_permits);
let put_object_min_size_bytes = rustfs_utils::get_env_usize(
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
@@ -327,7 +340,8 @@ impl ForegroundWriteAdmissionPolicy {
);
Self::Large {
gate: ForegroundWriteAdmissionGate::new(large_limit, wait_timeout),
gate: ForegroundWriteAdmissionGate::new(permit_limit, wait_timeout),
large_request_permits: u32::from(large_request_permits),
put_object_min_size_bytes,
multipart_part_min_size_bytes,
multipart_wait_timeout,
@@ -375,6 +389,7 @@ impl ForegroundWriteAdmissionPolicy {
if enabled && limit > 0 {
Self::Large {
gate: ForegroundWriteAdmissionGate::new(limit, wait_timeout),
large_request_permits: 1,
put_object_min_size_bytes: min_size_bytes,
multipart_part_min_size_bytes: 0,
multipart_wait_timeout,
@@ -392,21 +407,27 @@ impl ForegroundWriteAdmissionPolicy {
) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
match self {
Self::Disabled | Self::LegacyCounterOnly => Ok(ForegroundWriteAdmission::Disabled),
Self::Strict(gate) => gate.admit().await,
Self::Strict(gate) => gate.admit(1).await,
Self::Large {
gate,
large_request_permits,
put_object_min_size_bytes,
multipart_part_min_size_bytes,
multipart_wait_timeout,
multipart_max_pending,
} => match kind {
ForegroundWriteAdmissionKind::PutObject if should_gate_foreground_write(size, *put_object_min_size_bytes) => {
gate.admit().await
gate.admit(*large_request_permits).await
}
ForegroundWriteAdmissionKind::MultipartPart
if should_gate_foreground_write(size, *multipart_part_min_size_bytes) =>
{
gate.admit_queued(*multipart_wait_timeout, *multipart_max_pending).await
gate.admit_queued(
*multipart_wait_timeout,
*multipart_max_pending,
multipart_admission_permits(size, *large_request_permits),
)
.await
}
_ => Ok(ForegroundWriteAdmission::Disabled),
},
@@ -463,6 +484,15 @@ fn derive_multipart_admission_max_pending(configured_max_pending: usize, limit:
limit.saturating_mul(DERIVED_MULTIPART_ADMISSION_MAX_PENDING_FACTOR)
}
fn multipart_admission_permits(size: i64, large_request_permits: u32) -> u32 {
let Ok(size) = u64::try_from(size) else {
return large_request_permits;
};
u32::try_from(size.div_ceil(MULTIPART_ADMISSION_UNIT_BYTES))
.unwrap_or(large_request_permits)
.clamp(1, large_request_permits)
}
fn derive_large_put_admission_limit(configured_limit: usize, max_disk_reads: usize) -> usize {
if configured_limit > 0 {
return configured_limit;
@@ -1801,6 +1831,282 @@ mod integration_tests {
});
}
fn automatic_write_policy(configured_limit: Option<&str>, disk_read_limit: usize) -> ForegroundWriteAdmissionPolicy {
temp_env::with_vars(
[
(rustfs_config::ENV_PUT_FOREGROUND_ADMISSION_ENABLE, Some("false")),
(rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE, Some("true")),
(rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT, configured_limit),
(rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES, None),
(rustfs_config::ENV_PUT_MULTIPART_FOREGROUND_ADMISSION_MIN_SIZE_BYTES, None),
(rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS, Some("0")),
(rustfs_config::ENV_PUT_MULTIPART_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS, None),
(rustfs_config::ENV_PUT_MULTIPART_FOREGROUND_ADMISSION_MAX_PENDING, None),
],
|| ForegroundWriteAdmissionPolicy::from_env(disk_read_limit),
)
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_completes_bounded_multipart_bursts() {
use super::ForegroundWriteAdmissionKind::MultipartPart;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
for concurrency in [128, 256, 384] {
let policy = Arc::new(automatic_write_policy(None, 64));
let active = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let mut tasks = tokio::task::JoinSet::new();
for _ in 0..concurrency {
let (policy, active, peak) = (policy.clone(), active.clone(), peak.clone());
tasks.spawn(async move {
let admission = policy.admit(MultipartPart, 8 * 1024 * 1024).await.expect("gate stays open");
let ForegroundWriteAdmission::Admitted(permit) = admission else {
return false;
};
peak.fetch_max(active.fetch_add(1, Ordering::SeqCst) + 1, Ordering::SeqCst);
// A bounded transfer may occupy a permit for several seconds.
tokio::time::sleep(Duration::from_secs(3)).await;
active.fetch_sub(1, Ordering::SeqCst);
drop(permit);
true
});
}
while let Some(result) = tasks.join_next().await {
assert!(result.expect("upload task joins"), "{concurrency}-part burst must finish without a retry");
}
assert!(peak.load(Ordering::SeqCst) <= 256, "small parts retain a hard concurrency bound");
assert_eq!(policy.snapshot(64).active, Some(0));
assert_eq!(policy.snapshot(64).queued, Some(0));
}
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_preserves_large_and_unknown_write_bound() {
use super::ForegroundWriteAdmissionKind::{MultipartPart, PutObject};
for (kind, size) in [
(PutObject, 1024 * 1024 * 1024),
(MultipartPart, -1),
(MultipartPart, 64 * 1024 * 1024),
] {
let policy = automatic_write_policy(None, 64);
let mut held = Vec::new();
for _ in 0..32 {
let admission = policy.admit(kind, size).await.expect("gate stays open");
assert!(matches!(admission, ForegroundWriteAdmission::Admitted(_)));
held.push(admission);
}
assert!(matches!(
policy.admit(kind, size).await.expect("gate stays open"),
ForegroundWriteAdmission::Rejected
));
assert_eq!(policy.snapshot(64).queued, Some(0));
drop(held);
assert_eq!(policy.snapshot(64).active, Some(0));
}
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_shares_budget_between_write_sizes() {
use super::ForegroundWriteAdmissionKind::{MultipartPart, PutObject};
let policy = automatic_write_policy(None, 64);
let mut large = Vec::new();
for _ in 0..31 {
large.push(
policy
.admit(PutObject, 1024 * 1024 * 1024)
.await
.expect("large write admitted"),
);
}
let mut small = Vec::new();
for _ in 0..8 {
let admission = policy.admit(MultipartPart, 8 * 1024 * 1024).await.expect("gate stays open");
assert!(matches!(admission, ForegroundWriteAdmission::Admitted(_)));
small.push(admission);
}
assert!(matches!(
policy.admit(MultipartPart, 1).await.expect("gate stays open"),
ForegroundWriteAdmission::Rejected
));
assert!(matches!(
policy.admit(PutObject, 1024 * 1024 * 1024).await.expect("gate stays open"),
ForegroundWriteAdmission::Rejected
));
drop(small);
let last_large = policy
.admit(PutObject, 1024 * 1024 * 1024)
.await
.expect("released units are reusable");
assert!(matches!(last_large, ForegroundWriteAdmission::Admitted(_)));
drop((last_large, large));
assert_eq!(policy.snapshot(64).active, Some(0));
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_keeps_explicit_request_limit() {
use super::ForegroundWriteAdmissionKind::{MultipartPart, PutObject};
let policy = automatic_write_policy(Some("2"), 64);
let first = policy.admit(MultipartPart, 1).await.expect("first small part admitted");
let second = policy
.admit(PutObject, 1024 * 1024 * 1024)
.await
.expect("large write admitted");
assert!(matches!(first, ForegroundWriteAdmission::Admitted(_)));
assert!(matches!(second, ForegroundWriteAdmission::Admitted(_)));
assert_eq!(policy.snapshot(64).limit, Some(2));
assert!(matches!(
policy.admit(MultipartPart, 1).await.expect("gate stays open"),
ForegroundWriteAdmission::Rejected
));
drop((first, second));
assert_eq!(policy.snapshot(64).active, Some(0));
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_cancellation_releases_partial_reservation() {
use super::ForegroundWriteAdmissionKind::MultipartPart;
use std::sync::Arc;
let policy = Arc::new(automatic_write_policy(None, 1));
let held = policy
.admit(MultipartPart, 8 * 1024 * 1024)
.await
.expect("small part admitted");
let large_policy = policy.clone();
let large = tokio::spawn(async move { large_policy.admit(MultipartPart, -1).await });
tokio::task::yield_now().await;
assert_eq!(policy.snapshot(1).queued, Some(1));
let small_policy = policy.clone();
let small = tokio::spawn(async move { small_policy.admit(MultipartPart, 8 * 1024 * 1024).await });
tokio::task::yield_now().await;
assert!(!small.is_finished(), "a small part cannot steal a queued large request's reservation");
large.abort();
assert!(large.await.expect_err("large waiter cancelled").is_cancelled());
let resumed = small.await.expect("small waiter joins").expect("gate stays open");
assert!(matches!(resumed, ForegroundWriteAdmission::Admitted(_)));
assert_eq!(policy.snapshot(1).active, Some(2));
assert_eq!(policy.snapshot(1).queued, Some(0));
drop((resumed, held));
let whole = policy
.admit(MultipartPart, -1)
.await
.expect("all reserved units were returned");
assert!(matches!(whole, ForegroundWriteAdmission::Admitted(_)));
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_rounds_sizes_up_and_recovers_after_timeout() {
use super::ForegroundWriteAdmissionKind::MultipartPart;
for (size, admitted) in [
(0, 8),
(1, 8),
(8 * 1024 * 1024, 8),
(8 * 1024 * 1024 + 1, 4),
(16 * 1024 * 1024 + 1, 2),
(64 * 1024 * 1024, 1),
(i64::MAX, 1),
(-1, 1),
] {
let policy = automatic_write_policy(None, 1);
let mut held = Vec::new();
for _ in 0..admitted {
let admission = policy.admit(MultipartPart, size).await.expect("gate stays open");
assert!(matches!(admission, ForegroundWriteAdmission::Admitted(_)), "size {size}");
held.push(admission);
}
assert!(matches!(
policy.admit(MultipartPart, size).await.expect("gate stays open"),
ForegroundWriteAdmission::Rejected
));
assert_eq!(policy.snapshot(1).queued, Some(0));
drop(held);
assert_eq!(policy.snapshot(1).active, Some(0), "timeout must return any partial reservation");
assert!(matches!(
policy.admit(MultipartPart, -1).await.expect("all units reusable"),
ForegroundWriteAdmission::Admitted(_)
));
}
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_does_not_starve_queued_large_parts() {
use super::ForegroundWriteAdmissionKind::MultipartPart;
use std::sync::Arc;
let policy = Arc::new(automatic_write_policy(None, 1));
let held = policy.admit(MultipartPart, 1).await.expect("small part admitted");
let large_policy = policy.clone();
let large = tokio::spawn(async move { large_policy.admit(MultipartPart, -1).await });
tokio::task::yield_now().await;
let small_policy = policy.clone();
let small = tokio::spawn(async move { small_policy.admit(MultipartPart, 1).await });
tokio::task::yield_now().await;
assert_eq!(policy.snapshot(1).queued, Some(2));
drop(held);
let large_admission = large.await.expect("large waiter joins").expect("gate stays open");
assert!(matches!(large_admission, ForegroundWriteAdmission::Admitted(_)));
assert!(
!small.is_finished(),
"the large waiter owns the entire budget before the later small waiter"
);
drop(large_admission);
assert!(matches!(
small.await.expect("small waiter joins").expect("gate stays open"),
ForegroundWriteAdmission::Admitted(_)
));
}
#[tokio::test(start_paused = true)]
#[serial]
async fn test_concurrency_manager_auto_admission_keeps_queue_bound_in_requests() {
use super::ForegroundWriteAdmissionKind::MultipartPart;
use std::sync::Arc;
let policy = Arc::new(automatic_write_policy(None, 1));
let held = policy.admit(MultipartPart, -1).await.expect("large part fills the budget");
let started = tokio::time::Instant::now();
let mut tasks = tokio::task::JoinSet::new();
for _ in 0..17 {
let policy = policy.clone();
tasks.spawn(async move { policy.admit(MultipartPart, 1).await });
}
let overflow = tasks
.join_next()
.await
.expect("overflow waiter completes")
.expect("waiter joins")
.expect("gate stays open");
assert!(matches!(overflow, ForegroundWriteAdmission::Rejected));
assert_eq!(
tokio::time::Instant::now(),
started,
"a full queue rejects without waiting for its deadline"
);
assert_eq!(policy.snapshot(1).queued, Some(16), "unit subdivision does not expand the queue");
drop(held);
while let Some(result) = tasks.join_next().await {
assert!(matches!(
result.expect("waiter joins").expect("gate stays open"),
ForegroundWriteAdmission::Admitted(_)
));
}
assert_eq!(policy.snapshot(1).active, Some(0));
assert_eq!(policy.snapshot(1).queued, Some(0));
}
#[test]
fn test_concurrency_manager_derives_multipart_admission_max_pending_from_limit() {
assert_eq!(derive_multipart_admission_max_pending(7, 32), 7);
+271
View File
@@ -0,0 +1,271 @@
#!/usr/bin/env python3
"""Verify concurrent multipart uploads, per-attempt errors, and full readback.
Install boto3 and provide AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY for a disposable
test endpoint. Each invocation creates and removes its own uniquely named bucket.
Example: python3 scripts/verify_multipart_admission.py --endpoint http://localhost:9000 --uploads 24 --workers 16 --size-mib 400 --out /tmp/multipart-rung4
The issue #7385 reproduction used boto3/botocore 1.43.66 and urllib3 2.7.0.
By default each part gets one attempt, so throttling cannot be hidden by retries.
"""
import argparse
import concurrent.futures as cf
import hashlib
import json
import os
import platform
import threading
import time
import uuid
from collections import Counter
from pathlib import Path
import boto3
import botocore
import urllib3
from botocore.config import Config
from botocore.exceptions import BotoCoreError, ClientError
p = argparse.ArgumentParser()
p.add_argument("--endpoint", required=True)
p.add_argument("--uploads", type=int, default=16)
p.add_argument("--workers", type=int, default=16)
p.add_argument("--size-mib", type=int, default=400)
p.add_argument("--part-mib", type=int, default=8)
p.add_argument("--connect-timeout", type=int, default=30)
p.add_argument("--attempts", type=int, default=1)
p.add_argument("--out", required=True)
args = p.parse_args()
if any(
getattr(args, key) <= 0
for key in [
"uploads",
"workers",
"size_mib",
"part_mib",
"connect_timeout",
"attempts",
]
):
p.error("counts, sizes, timeouts, and attempts must be positive")
if args.part_mib < 5 and args.size_mib > args.part_mib:
p.error("non-final multipart parts must be at least 5 MiB")
if not os.environ.get("AWS_ACCESS_KEY_ID") or not os.environ.get(
"AWS_SECRET_ACCESS_KEY"
):
p.error("set AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY for the test endpoint")
out = Path(args.out)
out.mkdir(parents=True, exist_ok=False)
os.environ["NO_PROXY"] = "*"
config = Config(
connect_timeout=args.connect_timeout,
read_timeout=120,
retries={"total_max_attempts": args.attempts, "mode": "standard"},
max_pool_connections=args.uploads * args.workers,
s3={"addressing_style": "path"},
proxies={},
)
client = boto3.client(
"s3",
endpoint_url=args.endpoint,
region_name="us-east-1",
aws_access_key_id=os.environ["AWS_ACCESS_KEY_ID"],
aws_secret_access_key=os.environ["AWS_SECRET_ACCESS_KEY"],
aws_session_token=os.environ.get("AWS_SESSION_TOKEN"),
config=config,
)
bucket = "multipart-admission-" + uuid.uuid4().hex[:20]
size = args.size_mib * 1024 * 1024
part_size = args.part_mib * 1024 * 1024
payload = os.urandom(part_size)
part_count = (size + part_size - 1) // part_size
expected = hashlib.sha256()
for part_index in range(part_count):
expected.update(payload[: min(part_size, size - part_index * part_size)])
expected_digest = expected.hexdigest()
uploads = []
records, attempts = [], []
lock = threading.Lock()
start = time.monotonic()
def retry_event(**kw):
response = kw.get("response")
parsed = response[1] if response else {}
request = kw.get("request_dict") or {}
with lock:
attempts.append(
{
"t": time.monotonic() - start,
"attempt": kw.get("attempts"),
"query": request.get("query_string"),
"status": parsed.get("ResponseMetadata", {}).get("HTTPStatusCode"),
"error": parsed.get("Error"),
"exception": repr(kw.get("caught_exception")),
}
)
client.meta.events.register("needs-retry.s3.UploadPart", retry_event)
barrier = threading.Barrier(args.uploads)
def upload_part(key, upload_id, part_number):
began = time.monotonic()
record = {"key": key, "part": part_number, "start": began - start}
try:
body = payload[: min(part_size, size - (part_number - 1) * part_size)]
response = client.upload_part(
Bucket=bucket,
Key=key,
UploadId=upload_id,
PartNumber=part_number,
Body=body,
ContentLength=len(body),
)
record.update(
ok=True, etag=response["ETag"], metadata=response["ResponseMetadata"]
)
except (BotoCoreError, ClientError, OSError) as exc:
record.update(
ok=False,
error_type=type(exc).__name__,
error=str(exc),
response=getattr(exc, "response", None),
inner=repr(getattr(exc, "kwargs", {}).get("error")),
)
record["duration"] = time.monotonic() - began
with lock:
records.append(record)
return record
def upload_object(item):
key, upload_id = item
barrier.wait()
with cf.ThreadPoolExecutor(max_workers=args.workers) as pool:
parts = list(
pool.map(
lambda part_number: upload_part(key, upload_id, part_number),
range(1, part_count + 1),
)
)
result = {"key": key, "ok": all(r["ok"] for r in parts)}
if result["ok"]:
try:
response = client.complete_multipart_upload(
Bucket=bucket,
Key=key,
UploadId=upload_id,
MultipartUpload={
"Parts": [
{"PartNumber": r["part"], "ETag": r["etag"]} for r in parts
]
},
)
result["metadata"] = response["ResponseMetadata"]
except (BotoCoreError, ClientError, OSError) as exc:
result.update(ok=False, error_type=type(exc).__name__, error=str(exc))
result["finished"] = time.monotonic() - start
print(json.dumps(result), flush=True)
return result
bucket_created = False
summary = None
cleanup_errors = []
try:
client.create_bucket(Bucket=bucket)
bucket_created = True
for i in range(args.uploads):
key = f"object-{i:02}"
upload_id = client.create_multipart_upload(Bucket=bucket, Key=key)["UploadId"]
uploads.append((key, upload_id))
started_at_unix = time.time()
start = time.monotonic()
with cf.ThreadPoolExecutor(max_workers=args.uploads) as pool:
objects = list(pool.map(upload_object, uploads))
elapsed = time.monotonic() - start
upload_finished_at_unix = time.time()
for result in objects:
if result["ok"]:
try:
response = client.get_object(Bucket=bucket, Key=result["key"])
digest, received = hashlib.sha256(), 0
stream = response["Body"]
try:
for chunk in stream.iter_chunks(1024 * 1024):
digest.update(chunk)
received += len(chunk)
finally:
stream.close()
result["readback_bytes"] = received
result["readback_sha256"] = digest.hexdigest()
result["verified"] = (
received == size and result["readback_sha256"] == expected_digest
)
if not result["verified"]:
result["verification_error"] = "body length or SHA-256 mismatch"
except (BotoCoreError, ClientError, OSError) as exc:
result.update(verified=False, verification_error=str(exc))
summary = {
"config": vars(args),
"boto3": boto3.__version__,
"botocore": botocore.__version__,
"urllib3": urllib3.__version__,
"platform": platform.platform(),
"effective_retries": client.meta.config.retries,
"elapsed": elapsed,
"started_at_unix": started_at_unix,
"upload_finished_at_unix": upload_finished_at_unix,
"expected_sha256": expected_digest,
"objects_passed": sum(r["ok"] for r in objects),
"objects_verified": sum(r.get("verified", False) for r in objects),
"parts_passed": sum(r["ok"] for r in records),
"parts_total": len(records),
"errors": dict(
Counter(
(r.get("response") or {})
.get("Error", {})
.get("Code", r.get("error_type"))
for r in records
if not r["ok"]
)
),
"attempts": len(attempts),
"objects": objects,
}
(out / "summary.json").write_text(json.dumps(summary, indent=2))
print(json.dumps(summary, indent=2), flush=True)
finally:
(out / "parts.json").write_text(json.dumps(records, indent=2))
(out / "attempts.json").write_text(json.dumps(attempts, indent=2))
for key, upload_id in uploads:
try:
client.abort_multipart_upload(Bucket=bucket, Key=key, UploadId=upload_id)
except ClientError as exc:
if exc.response.get("Error", {}).get("Code") != "NoSuchUpload":
cleanup_errors.append(
{"key": key, "operation": "abort", "error": str(exc)}
)
except (BotoCoreError, OSError) as exc:
cleanup_errors.append({"key": key, "operation": "abort", "error": str(exc)})
try:
client.delete_object(Bucket=bucket, Key=key)
except (BotoCoreError, ClientError, OSError) as exc:
cleanup_errors.append(
{"key": key, "operation": "delete", "error": str(exc)}
)
if bucket_created:
try:
client.delete_bucket(Bucket=bucket)
except (BotoCoreError, ClientError, OSError) as exc:
cleanup_errors.append({"bucket": bucket, "error": str(exc)})
if cleanup_errors:
(out / "cleanup-errors.json").write_text(json.dumps(cleanup_errors, indent=2))
client.close()
raise SystemExit(
0
if summary and summary["objects_verified"] == args.uploads and not cleanup_errors
else 1
)