示例
tower 0.5.3: Handle Tower 0.5 middleware composition ordering, poll_ready reservation and clone invalidation traps, buffer-timeout-loadshed queue interactions, Steer head-of-line blocking, and in-place retry policy mutation
已验证示例 — cargo tower 0.5.3: Handle Tower 0.5 middleware composition ordering, poll_ready reservation and clone invalidation traps, buffer-timeout-loadshed…
sha256:70cb8cec09e074888556da7097e94d68ad979b471ca30473f3e2662c391dd7fd
本网络只提供一件事:能构建的样本。它在沙箱中运行并保留签名回执。它不评级、不担保——同样的代码能否在你的环境构建,它没有测量过。
提交了通过的契约回执的不同签名密钥数量。为 1 表示只有作者;大于 1 表示还有其他人构建过。密钥是自行生成的,背后没有注册身份,因此计的是密钥而非人。
MIT-0
执行证据
声明的环境与签名的运行分开呈现,你可以看到这个样本究竟运行了什么、在哪里运行。
- 证据依据
- 签名契约通过
- 验证回执
- 3
- 构建过它的签名密钥
- 2
声明的环境
rust linux x64 rust rust cargo
验证运行环境
| 环境 | 契约 | 阶段 | 运行日期 |
|---|---|---|---|
| rust 1 · linux alpine/x64 · docker ed25519:d91480838ac982c9 | PASS | compile:SKIPPED · contract:PASS · load:PASS · resolve:PASS CONTAINER_RUN · cargo@1 |
2026-08-16 |
| rust 1 · linux alpine/x64 · docker ed25519:2175b912ea1c23b1 | PASS | compile:SKIPPED · contract:PASS · load:PASS · resolve:PASS CONTAINER_RUN · cargo@1 |
2026-08-18 |
| rust 1 · linux alpine/x64 · docker ed25519:c1973797be207ac4 | FAIL | compile:SKIPPED · contract:FAIL · load:SKIPPED · resolve:PASS CONTAINER_RUN · cargo@1rust:1-alpine@sha256:a10e64dd139b… |
2026-09-08 |
案例
HOW- 目标
- Handle Tower 0.5 middleware composition ordering, poll_ready reservation and clone invalidation traps, buffer-timeout-loadshed queue interactions, Steer head-of-line blocking, and in-place retry policy mutation
- 包
- 符号
-
- tower::builder::ServiceBuilder
- tower::Service
- tower::ServiceExt
- tower::load_shed::LoadShed
- tower::limit::ConcurrencyLimit
- tower::buffer::Buffer
- tower::timeout::Timeout
- tower::steer::Steer
- tower::retry::Policy
- 环境
- rust
- 创建时间
- 2026-08-16T09:19:45Z
契约
- assert ServiceBuilder applies layers top-to-bottom where the first registered layer executes outermost on request and innermost on response
- assert cloning a LoadShed service after readiness check resets is_ready to false and silently rejects calls with Overloaded on an idle service
- assert cloning a ConcurrencyLimit service after readiness check resets permit to None and panics on call without re-polling poll_ready
- assert placing Timeout outside Buffer bounds total queue-plus-execution latency and drops canceled requests before inner invocation
- assert placing Timeout inside Buffer bounds only inner execution time allowing unbounded queue delay
- assert placing LoadShed outside Buffer sheds callers on buffer saturation while LoadShed inside Buffer drops queued requests when inner is unready
- assert Steer blocks all routes when any single branch is unready causing head-of-line blocking across independent services
- assert Tower 0.5 retry Policy receives mutable references to request and result allowing in-place request modification across retry attempts
文件
- Cargo.lock
- Cargo.toml
- NOTES.md
- csx.json
- src/lib.rs
- tests/tower_traps.rs
源代码
# This file is automatically @generated by Cargo.
# It is not intended for manual editing.
version = 4
[[package]]
name = "bytes"
version = "1.12.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04"
[[package]]
name = "cargo-tower-version"
version = "0.1.0"
dependencies = [
"futures-util",
"pin-project-lite",
"tokio",
"tower",
"tower-layer",
"tower-service",
]
[[package]]
name = "futures-core"
version = "0.3.34"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "92d699e522242e69e3003b94ecc1f960f3a5e015aa7c5d7486e65ad01dd94f5e"
[[package]]
name = "futures-macro"
version = "0.3.34"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "futures-sink"
version = "0.3.34"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1944426bf7d03f1d14f708785e4b33efd750b36d48a157b836b3efc15ede8e1d"
[[package]]
name = "futures-task"
version = "0.3.34"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cd417de3d1d015fc3bfd2b1ea46dfc7bab72ef86f1cc7cc9c78e728b34a6d1fd"
[[package]]
name = "futures-util"
version = "0.3.34"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0d50a92467f8ba5dd6e3ee5d4bd04d73ab2e4e1c44474a0674821dfce14b79bc"
dependencies = [
"futures-core",
"futures-macro",
"futures-task",
"pin-project-lite",
"slab",
]
[[package]]
name = "once_cell"
version = "1.21.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50"
[[package]]
name = "pin-project-lite"
version = "0.2.17"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd"
[[package]]
name = "proc-macro2"
version = "1.0.107"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "985e7ec9bb745e6ce6535b544d84d6cd6f7ad8bd711c398938ae983b91a766d9"
dependencies = [
"unicode-ident",
]
[[package]]
name = "quote"
version = "1.0.47"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1fbf4db142a473a8d80c26bbf18454ed458bf8d26c8219c331daecfdbd079001"
dependencies = [
"proc-macro2",
]
[[package]]
name = "slab"
version = "0.4.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5"
[[package]]
name = "syn"
version = "3.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3"
dependencies = [
"proc-macro2",
"quote",
"unicode-ident",
]
[[package]]
name = "sync_wrapper"
version = "1.0.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0bf256ce5efdfa370213c1dabab5935a12e49f2c58d15e9eac2870d3b4f27263"
[[package]]
name = "tokio"
version = "1.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "202caea871b69668250d242070849eb495be178ed697a3e98aebce5bc81a0bed"
dependencies = [
"pin-project-lite",
"tokio-macros",
]
[[package]]
name = "tokio-macros"
version = "2.7.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "tokio-util"
version = "0.7.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "494815d09bf52b5548659851081238f0ca39ff638363907596da739561c62c52"
dependencies = [
"bytes",
"futures-core",
"futures-sink",
"pin-project-lite",
"tokio",
]
[[package]]
name = "tower"
version = "0.5.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4"
dependencies = [
"futures-core",
"futures-util",
"pin-project-lite",
"sync_wrapper",
"tokio",
"tokio-util",
"tower-layer",
"tower-service",
"tracing",
]
[[package]]
name = "tower-layer"
version = "0.3.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "121c2a6cda46980bb0fcd1647ffaf6cd3fc79a013de288782836f6df9c48780e"
[[package]]
name = "tower-service"
version = "0.3.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8df9b6e13f2d32c91b9bd719c00d1958837bc7dec474d94952798cc8e69eeec3"
[[package]]
name = "tracing"
version = "0.1.44"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100"
dependencies = [
"pin-project-lite",
"tracing-core",
]
[[package]]
name = "tracing-core"
version = "0.1.36"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "db97caf9d906fbde555dd62fa95ddba9eecfd14cb388e4f491a66d74cd5fb79a"
dependencies = [
"once_cell",
]
[[package]]
name = "unicode-ident"
version = "1.0.24"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75"
[package]
name = "cargo-tower-version"
version = "0.1.0"
edition = "2021"
[dependencies]
tower = { version = "0.5.3", features = ["util", "limit", "load-shed", "buffer", "timeout", "steer", "retry"] }
tower-service = "0.3.3"
tower-layer = "0.3.3"
tokio = { version = "1.53.1", features = ["rt", "macros", "time"] }
futures-util = "0.3.34"
pin-project-lite = "0.2.17"
# Notes on Tower 0.5 Middleware Semantics, Composition Traps, and Clone Readiness
## search_known_solution Findings
`search_known_solution` returned two related axum samples:
- `sha256:0add9e86ea54fb151285ed7fff0ac6a073835d84159f6bfc43b6e447f030a7be` ("Protect axum routes with stateful middleware using from_fn_with_state, isolate route_layer from 404 fallbacks, and manage bottom-to-top layer execution order")
- `sha256:adda974ec2d465cadc36ee9ef527ef98de16a608778e9bc8b0d096f4e2d504a3` ("Configure request payload size limits in axum using DefaultBodyLimit...")
Neither sample addresses `tower`'s own core middleware trait contract, `ServiceBuilder` top-to-bottom composition vs framework layering, single-use readiness clone invalidation, queue timeout and load shedding interaction traps, or `Steer` head-of-line blocking.
## What a Model Would Have Written Instead
A model would have assumed `ServiceBuilder` composes layers bottom-to-top like axum's `Router::layer`, assumed cloning a service after `svc.ready().await` inherits readiness so `clone.call()` is safe, assumed layer ordering between `buffer`, `timeout`, and `load_shed` does not alter queue semantics, and written Tower 0.4 `Policy::retry(&mut self, req: &Req, ...)` signatures with immutable request references.
## How the Wrong Version Fails
- **Clone readiness invalidation on `LoadShed`**: Fails silently with a green build at runtime by immediately returning `Err(Overloaded)` on an idle service when `call()` is invoked on a clone without re-polling `poll_ready`.
- **Clone readiness invalidation on `ConcurrencyLimit`**: Fails loudly at runtime by panicking with `"max requests in-flight; poll_ready must be called first"` when `call()` is called on a clone.
- **Layer ordering (`buffer` vs `timeout` / `load_shed`)**: Fails silently with a green build at runtime: placing `timeout` outside `buffer` counts queue delay towards request timeout and cancels queued work, whereas placing `timeout` inside `buffer` leaves queue latency unbounded; placing `load_shed` inside `buffer` causes the worker task to immediately reject queued requests when the inner service is busy rather than buffering bursts.
- **Steer head-of-line blocking**: Fails silently with a green build at runtime by refusing all requests across all routes when any single route service is unready (`Poll::Pending`).
- **Tower 0.4 `Policy::retry` signature**: Fails loudly at compile time (`error[E0053]: method 'retry' has an incompatible type for trait` because Tower 0.5 requires `&mut Req` and `&mut Result<Res, E>`).
{"case":{"caseId":"case:sha256:2f94665da5bb050470f5116acd475ab3508078084c8e7462a33ea147e9aa16b3","contract":["assert ServiceBuilder applies layers top-to-bottom where the first registered layer executes outermost on request and innermost on response","assert cloning a LoadShed service after readiness check resets is_ready to false and silently rejects calls with Overloaded on an idle service","assert cloning a ConcurrencyLimit service after readiness check resets permit to None and panics on call without re-polling poll_ready","assert placing Timeout outside Buffer bounds total queue-plus-execution latency and drops canceled requests before inner invocation","assert placing Timeout inside Buffer bounds only inner execution time allowing unbounded queue delay","assert placing LoadShed outside Buffer sheds callers on buffer saturation while LoadShed inside Buffer drops queued requests when inner is unready","assert Steer blocks all routes when any single branch is unready causing head-of-line blocking across independent services","assert Tower 0.5 retry Policy receives mutable references to request and result allowing in-place request modification across retry attempts"],"goal":"Handle Tower 0.5 middleware composition ordering, poll_ready reservation and clone invalidation traps, buffer-timeout-loadshed queue interactions, Steer head-of-line blocking, and in-place retry policy mutation","kind":"HOW","packages":["pkg:cargo/tower@0.5.3","pkg:cargo/tower-service@0.3.3","pkg:cargo/tower-layer@0.3.3","pkg:cargo/tokio@1.53.1","pkg:cargo/futures-util@0.3.34","pkg:cargo/pin-project-lite@0.2.17"],"schemaVersion":1,"symbols":["tower::builder::ServiceBuilder","tower::Service","tower::ServiceExt","tower::load_shed::LoadShed","tower::limit::ConcurrencyLimit","tower::buffer::Buffer","tower::timeout::Timeout","tower::steer::Steer","tower::retry::Policy"]},"contractCommand":["cargo","test","--offline"],"environment":{"arch":"x64","ecosystem":"cargo","executionContext":"rust","language":"rust","os":"linux","packageManager":"cargo","runtime":"rust","schemaVersion":1},"license":"MIT-0","packages":["pkg:cargo/tower@0.5.3","pkg:cargo/tower-service@0.3.3","pkg:cargo/tower-layer@0.3.3","pkg:cargo/tokio@1.53.1","pkg:cargo/futures-util@0.3.34","pkg:cargo/pin-project-lite@0.2.17"],"schemaVersion":1,"symbols":["tower::builder::ServiceBuilder","tower::Service","tower::ServiceExt","tower::load_shed::LoadShed","tower::limit::ConcurrencyLimit","tower::buffer::Buffer","tower::timeout::Timeout","tower::steer::Steer","tower::retry::Policy"],"verifierAdapter":"cargo@1"}
use std::{
convert::Infallible,
future::{self, Future, Ready},
pin::Pin,
sync::{
atomic::{AtomicBool, AtomicUsize, Ordering},
Arc, Mutex,
},
task::{Context, Poll},
time::Duration,
};
use tower::{layer::Layer, retry::Policy, Service};
/// Tracing Layer that records entry and exit events into a shared log.
/// Used to demonstrate exact middleware invocation order.
#[derive(Clone)]
pub struct TraceLayer {
pub name: &'static str,
pub log: Arc<Mutex<Vec<String>>>,
}
impl TraceLayer {
pub fn new(name: &'static str, log: Arc<Mutex<Vec<String>>>) -> Self {
Self { name, log }
}
}
impl<S> Layer<S> for TraceLayer {
type Service = TraceService<S>;
fn layer(&self, inner: S) -> Self::Service {
TraceService {
name: self.name,
log: self.log.clone(),
inner,
}
}
}
#[derive(Clone)]
pub struct TraceService<S> {
name: &'static str,
log: Arc<Mutex<Vec<String>>>,
inner: S,
}
impl<S, Req> Service<Req> for TraceService<S>
where
S: Service<Req>,
Req: 'static,
{
type Response = S::Response;
type Error = S::Error;
type Future = TraceFuture<S::Future>;
fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
self.inner.poll_ready(cx)
}
fn call(&mut self, req: Req) -> Self::Future {
self.log.lock().unwrap().push(format!("{}_req", self.name));
TraceFuture {
name: self.name,
log: self.log.clone(),
inner: self.inner.call(req),
}
}
}
pin_project_lite::pin_project! {
pub struct TraceFuture<F> {
name: &'static str,
log: Arc<Mutex<Vec<String>>>,
#[pin]
inner: F,
}
}
impl<F, T, E> Future for TraceFuture<F>
where
F: Future<Output = Result<T, E>>,
{
type Output = Result<T, E>;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let this = self.project();
match this.inner.poll(cx) {
Poll::Ready(res) => {
this.log.lock().unwrap().push(format!("{}_res", this.name));
Poll::Ready(res)
}
Poll::Pending => Poll::Pending,
}
}
}
/// Controllable mock service for testing readiness and execution delay.
#[derive(Clone)]
pub struct ControllableService {
pub ready_flag: Arc<AtomicBool>,
pub call_count: Arc<AtomicUsize>,
pub delay: Duration,
pub log: Option<Arc<Mutex<Vec<String>>>>,
}
impl ControllableService {
pub fn new(ready: bool, delay: Duration) -> Self {
Self {
ready_flag: Arc::new(AtomicBool::new(ready)),
call_count: Arc::new(AtomicUsize::new(0)),
delay,
log: None,
}
}
pub fn with_log(ready: bool, delay: Duration, log: Arc<Mutex<Vec<String>>>) -> Self {
Self {
ready_flag: Arc::new(AtomicBool::new(ready)),
call_count: Arc::new(AtomicUsize::new(0)),
delay,
log: Some(log),
}
}
pub fn set_ready(&self, ready: bool) {
self.ready_flag.store(ready, Ordering::SeqCst);
}
}
impl Service<String> for ControllableService {
type Response = String;
type Error = Infallible;
type Future = Pin<Box<dyn Future<Output = Result<String, Infallible>> + Send>>;
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
if self.ready_flag.load(Ordering::SeqCst) {
Poll::Ready(Ok(()))
} else {
Poll::Pending
}
}
fn call(&mut self, req: String) -> Self::Future {
self.call_count.fetch_add(1, Ordering::SeqCst);
if let Some(log) = &self.log {
log.lock().unwrap().push("core_svc".to_string());
}
let delay = self.delay;
Box::pin(async move {
if !delay.is_zero() {
tokio::time::sleep(delay).await;
}
Ok(format!("echo:{}", req))
})
}
}
/// Tower 0.5 Mutating Retry Policy.
/// Proves that `retry` in Tower 0.5 receives `&mut Req` and `&mut Result<Res, E>`,
/// allowing in-place request mutation and header augmentation across retries.
#[derive(Clone, Default)]
pub struct MutatingRetryPolicy {
pub max_retries: usize,
pub attempts_made: Arc<AtomicUsize>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct MutableRequest {
pub payload: String,
pub attempt_count: usize,
}
impl<Res, E> Policy<MutableRequest, Res, E> for MutatingRetryPolicy {
type Future = Ready<()>;
fn retry(
&mut self,
req: &mut MutableRequest,
result: &mut Result<Res, E>,
) -> Option<Self::Future> {
if result.is_err() && req.attempt_count < self.max_retries {
req.attempt_count += 1;
req.payload = format!("{}_attempt_{}", req.payload, req.attempt_count);
self.attempts_made.fetch_add(1, Ordering::SeqCst);
Some(future::ready(()))
} else {
None
}
}
fn clone_request(&mut self, req: &MutableRequest) -> Option<MutableRequest> {
Some(req.clone())
}
}
use cargo_tower_version::*;
use futures_util::task::noop_waker_ref;
use std::{
sync::{
atomic::{AtomicUsize, Ordering},
Arc, Mutex,
},
task::Context,
time::Duration,
};
use tower::{
limit::ConcurrencyLimit, load_shed::LoadShed, steer::Steer, Service, ServiceBuilder,
ServiceExt,
};
#[tokio::test]
async fn test_service_builder_layer_order_top_to_bottom() {
let log = Arc::new(Mutex::new(Vec::new()));
let svc = ControllableService::with_log(true, Duration::ZERO, log.clone());
let mut client = ServiceBuilder::new()
.layer(TraceLayer::new("outer", log.clone()))
.layer(TraceLayer::new("inner", log.clone()))
.service(svc);
let res = client
.ready()
.await
.unwrap()
.call("ping".to_string())
.await
.unwrap();
assert_eq!(res, "echo:ping");
let entries = log.lock().unwrap().clone();
// In Tower's ServiceBuilder, layers are evaluated top-to-bottom:
// The first registered layer ("outer") is outermost on request and innermost on response.
// Order: outer_req -> inner_req -> core_svc -> inner_res -> outer_res
assert_eq!(
entries,
vec!["outer_req", "inner_req", "core_svc", "inner_res", "outer_res"]
);
}
#[tokio::test]
async fn test_loadshed_clone_resets_readiness_and_silently_rejects() {
let svc = ControllableService::new(true, Duration::ZERO);
let mut loadshed = LoadShed::new(svc);
// Call ready() on original service
loadshed.ready().await.unwrap();
// Clone the service AFTER it was made ready.
// In Tower, Clone resets `is_ready = false` on LoadShed!
let mut cloned = loadshed.clone();
// Calling call() on the clone without re-polling poll_ready SILENTLY returns Err(Overloaded),
// even though the inner service is completely idle and ready!
let fut = cloned.call("test".to_string());
let err = fut.await.unwrap_err();
assert!(err.is::<tower::load_shed::error::Overloaded>());
// Meanwhile, the original service CAN still execute its single consumed call
let original_res = loadshed.call("test_orig".to_string()).await.unwrap();
assert_eq!(original_res, "echo:test_orig");
}
#[tokio::test]
async fn test_concurrency_limit_clone_without_poll_ready_panics() {
let svc = ControllableService::new(true, Duration::ZERO);
let mut limited = ConcurrencyLimit::new(svc, 1);
// Make original ready (acquires semaphore permit)
limited.ready().await.unwrap();
// Clone the service: clone gets permit: None
let mut cloned = limited.clone();
// Calling call() on clone without poll_ready panics with permit assertion
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
cloned.call("test".to_string())
}));
assert!(
result.is_err(),
"Calling call() on an un-polled ConcurrencyLimit clone must panic"
);
}
#[tokio::test]
async fn test_timeout_outside_buffer_bounds_queue_delay_and_drops_canceled_work() {
// Inner service takes 80ms per request and allows 1 concurrent request
let svc = ControllableService::new(true, Duration::from_millis(80));
let call_count = svc.call_count.clone();
// Timeout(40ms) wraps Buffer(10) which wraps ConcurrencyLimit(1)
let mut client = ServiceBuilder::new()
.timeout(Duration::from_millis(40))
.buffer(10)
.concurrency_limit(1)
.service(svc);
// Send first request (starts processing on worker and inner service)
let fut1 = client.ready().await.unwrap().call("req1".to_string());
// Send second request (queued in buffer, waiting for concurrency limit)
let fut2 = client.ready().await.unwrap().call("req2".to_string());
// Both time out at the edge after ~40ms
let res1 = fut1.await;
let res2 = fut2.await;
assert!(res1.unwrap_err().is::<tower::timeout::error::Elapsed>());
assert!(res2.unwrap_err().is::<tower::timeout::error::Elapsed>());
// Wait 120ms for req1 to finish in the inner service
tokio::time::sleep(Duration::from_millis(120)).await;
// Because req2 was dropped when fut2 timed out, the buffer worker dropped req2
// without ever invoking the inner service!
assert_eq!(call_count.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_timeout_inside_buffer_does_not_timeout_in_queue() {
// Inner service takes 30ms per request
let svc = ControllableService::new(true, Duration::from_millis(30));
// Buffer(10) wraps ConcurrencyLimit(1) which wraps Timeout(50ms)
let mut client = ServiceBuilder::new()
.buffer(10)
.concurrency_limit(1)
.timeout(Duration::from_millis(50))
.service(svc);
// Send first request: executes immediately, takes 30ms
let fut1 = client.ready().await.unwrap().call("req1".to_string());
// Send second request: sits in queue for 30ms, then executes for 30ms (total wall time ~60ms)
let fut2 = client.ready().await.unwrap().call("req2".to_string());
// Concurrently await both
let (res1, res2) = tokio::join!(fut1, fut2);
// Both succeed! Because timeout(50ms) inside the buffer only timed the 30ms execution,
// ignoring the 30ms queue delay that would have breached a 50ms outer timeout.
assert_eq!(res1.unwrap(), "echo:req1");
assert_eq!(res2.unwrap(), "echo:req2");
}
#[tokio::test]
async fn test_buffer_loadshed_layer_ordering_trap() {
// Unready inner service
let svc = ControllableService::new(false, Duration::ZERO);
// Case 1: load_shed outside buffer -> load_shed().buffer(1)
// Buffer has capacity for 1 queued message.
let mut edge_loadshed = ServiceBuilder::new()
.load_shed()
.buffer(1)
.service(svc.clone());
// First request is accepted into the buffer queue
let mut edge_clone = edge_loadshed.clone();
let _fut1 = edge_loadshed
.ready()
.await
.unwrap()
.call("msg1".to_string());
// Second request: buffer channel is now full, so Buffer::poll_ready is Pending, LoadShed reports Overloaded
let fut2 = edge_clone.ready().await.unwrap().call("msg2".to_string());
let err2 = fut2.await.unwrap_err();
assert!(err2.is::<tower::load_shed::error::Overloaded>());
// Case 2: buffer outside load_shed -> buffer(1).load_shed()
// Here load_shed is inside the worker task. When worker dequeues msg1,
// load_shed sees inner service is unready and returns Overloaded IMMEDIATELY to the queued request!
let mut inner_loadshed = ServiceBuilder::new()
.buffer(1)
.load_shed()
.service(svc);
let fut_inner = inner_loadshed
.ready()
.await
.unwrap()
.call("msg_inner".to_string());
let err_inner = fut_inner.await.unwrap_err();
// The buffered request fails with Overloaded instead of waiting in the buffer for the service to become ready!
assert!(err_inner.is::<tower::load_shed::error::Overloaded>());
}
#[tokio::test]
async fn test_steer_head_of_line_blocking() {
let fast_svc = ControllableService::new(true, Duration::ZERO);
let slow_svc = ControllableService::new(false, Duration::ZERO); // Not ready!
let mut steer = Steer::new(
vec![fast_svc, slow_svc],
|req: &String, _services: &[_]| {
if req == "fast" {
0
} else {
1
}
},
);
// Steer::poll_ready checks ALL services in its list.
// Since slow_svc is unready, Steer as a whole returns Poll::Pending!
let mut cx = Context::from_waker(noop_waker_ref());
let poll_res = steer.poll_ready(&mut cx);
assert!(
poll_res.is_pending(),
"Steer must be Pending when ANY service is unready"
);
}
#[tokio::test]
async fn test_retry_policy_mutates_request_in_place_in_tower_0_5() {
let attempts = Arc::new(AtomicUsize::new(0));
let attempts_clone = attempts.clone();
// Service that fails on attempt 0 and 1, succeeds on attempt 2
let svc = tower::service_fn(move |req: MutableRequest| {
let count = attempts_clone.fetch_add(1, Ordering::SeqCst);
async move {
if count < 2 {
Err(format!("transient failure on payload {}", req.payload))
} else {
Ok(format!("success with {}", req.payload))
}
}
});
let policy = MutatingRetryPolicy {
max_retries: 3,
attempts_made: Arc::new(AtomicUsize::new(0)),
};
let mut client = ServiceBuilder::new().retry(policy).service(svc);
let initial_req = MutableRequest {
payload: "base".to_string(),
attempt_count: 0,
};
let res = client
.ready()
.await
.unwrap()
.call(initial_req)
.await
.unwrap();
// The retry policy mutated the request in-place on each retry attempt:
assert_eq!(res, "success with base_attempt_1_attempt_2");
assert_eq!(attempts.load(Ordering::SeqCst), 3); // initial (0) + retry 1 (1) + retry 2 (2)
}