use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
use tropel_core::config::{
ExecutionConfig, HttpConfig, JobConfig, OutputConfig, Stage, ThinkTimeConfig, ThresholdConfig,
};
use tropel_engine::Engine;
use tropel_ext::registry::ExtensionRegistry;
use tropel_metrics::thresholds::evaluate_thresholds;
use tropel_sdk::Result;
async fn start_peak_server() -> (std::net::SocketAddr, Arc<AtomicUsize>) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let active = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let peak_out = peak.clone();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
let active = active.clone();
let peak = peak.clone();
tokio::spawn(async move {
let now = active.fetch_add(1, Ordering::SeqCst) + 1;
peak.fetch_max(now, Ordering::SeqCst);
let mut head = Vec::new();
let mut buf = [0u8; 4096];
loop {
let n = match sock.read(&mut buf).await {
Ok(0) | Err(_) => break,
Ok(n) => n,
};
head.extend_from_slice(&buf[..n]);
if head.windows(4).any(|w| w == b"\r\n\r\n") {
break;
}
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
let body = r#"{"ok":true}"#;
let resp = format!(
"HTTP/1.1 200 OK\r\nConnection: close\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
);
let _ = sock.write_all(resp.as_bytes()).await;
active.fetch_sub(1, Ordering::SeqCst);
});
}
});
(addr, peak_out)
}
async fn start_echo_server() -> std::net::SocketAddr {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
tokio::spawn(async move {
let mut buf = [0u8; 4096];
loop {
let mut head = Vec::new();
loop {
let n = match sock.read(&mut buf).await {
Ok(0) | Err(_) => return,
Ok(n) => n,
};
head.extend_from_slice(&buf[..n]);
if head.windows(4).any(|w| w == b"\r\n\r\n") {
break;
}
}
let body = r#"{"ok":true}"#;
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
);
if sock.write_all(resp.as_bytes()).await.is_err() {
return;
}
}
});
}
});
addr
}
async fn start_500_server() -> std::net::SocketAddr {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
tokio::spawn(async move {
let mut buf = [0u8; 4096];
loop {
let mut head = Vec::new();
loop {
let n = match sock.read(&mut buf).await {
Ok(0) | Err(_) => return,
Ok(n) => n,
};
head.extend_from_slice(&buf[..n]);
if head.windows(4).any(|w| w == b"\r\n\r\n") {
break;
}
}
let body = r#"{"error":"boom"}"#;
let resp = format!(
"HTTP/1.1 500 Internal Server Error\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
);
if sock.write_all(resp.as_bytes()).await.is_err() {
return;
}
}
});
}
});
addr
}
async fn start_hung_server() -> std::net::SocketAddr {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
tokio::spawn(async move {
let mut buf = [0u8; 4096];
loop {
match sock.read(&mut buf).await {
Ok(0) | Err(_) => break,
Ok(_) => {}
}
}
});
}
});
addr
}
fn write_k6_script(base: &str, tag: &str) -> String {
let dir = std::env::temp_dir();
let path = dir.join(format!("tropel-k6-e2e-{}-{}.js", std::process::id(), tag));
let script = format!(
r#"import http from 'k6/http';
import {{ check }} from 'k6';
export default function () {{
const res = http.get('{base}/');
check(res, {{ 'status is 200': (r) => r.status === 200 }});
}}
"#
);
std::fs::write(&path, script).unwrap();
path.to_string_lossy().to_string()
}
fn write_k6_timeout_script(base: &str, tag: &str, timeout: &str) -> String {
let dir = std::env::temp_dir();
let path = dir.join(format!("tropel-k6-e2e-{}-{}.js", std::process::id(), tag));
let script = format!(
r#"import http from 'k6/http';
import {{ check }} from 'k6';
export default function () {{
const res = http.get('{base}/', {{ timeout: '{timeout}' }});
check(res, {{ 'status is 200': (r) => r.status === 200 }});
}}
"#
);
std::fs::write(&path, script).unwrap();
path.to_string_lossy().to_string()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn k6_script_records_requests_checks_and_real_latency() -> Result<()> {
let srv = start_echo_server().await;
let script = write_k6_script(&format!("http://{srv}"), "check");
let mut thresholds = HashMap::new();
thresholds.insert(
"http_req_duration".to_string(),
ThresholdConfig {
expression: "http_req_duration.p95 < 5000".to_string(),
abort_on_fail: false,
delay_abort_eval: None,
},
);
let config = JobConfig {
input: script.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::ConstantVus {
vus: 2,
duration: "3s".to_string(),
graceful_stop: Some("5s".to_string()),
think_time: ThinkTimeConfig::default(),
},
thresholds,
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let result = engine.run(&config).await?;
let m = &result.metrics;
assert!(m.http_reqs > 0, "http_reqs > 0, got {}", m.http_reqs);
assert!(
m.checks_total > 0,
"checks_total > 0, got {}",
m.checks_total
);
assert_eq!(
m.checks_failed, 0,
"all checks passed, got {} failed of {} total",
m.checks_failed, m.checks_total
);
let dur = m
.http_req_duration
.as_ref()
.expect("http_req_duration summary present");
assert!(dur.count > 0, "http_req_duration has samples");
assert!(
dur.max > 0.0,
"http_req_duration max > 0 (real latency measured)"
);
let threshold_results = evaluate_thresholds(&result.effective_thresholds, m);
let t = threshold_results
.iter()
.find(|t| t.name == "http_req_duration")
.expect("threshold evaluated");
assert!(
t.passed,
"threshold '{}' passed (actual={} < {})",
t.expression, t.actual, t.threshold
);
let _ = std::fs::remove_file(&script);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn ramping_stages_span_wall_clock_and_reach_target() -> Result<()> {
let (srv, peak) = start_peak_server().await;
let coll = write_k6_script(&format!("http://{srv}"), "ramp");
let config = JobConfig {
input: coll.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::RampingVus {
start_vus: 1,
stages: vec![
Stage {
duration: "1s".to_string(),
target: 3,
},
Stage {
duration: "1s".to_string(),
target: 3,
},
Stage {
duration: "1s".to_string(),
target: 1,
},
],
graceful_ramp_down: Some("5s".to_string()),
graceful_stop: Some("5s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let start = Instant::now();
let result = engine.run(&config).await?;
let elapsed = start.elapsed();
let m = &result.metrics;
assert!(
elapsed >= std::time::Duration::from_millis(2250),
"ramping run elapsed {elapsed:?}, expected >= 2.25s (stages not collapsed)"
);
let peak_obs = peak.load(Ordering::SeqCst);
assert!(
peak_obs >= 3,
"peak concurrent connections >= 3, got {} (pool actually grew)",
peak_obs
);
assert!(
m.http_reqs > 0,
"http_reqs > 0 during ramp, got {}",
m.http_reqs
);
let _ = std::fs::remove_file(&coll);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn non_2xx_drives_http_req_failed() -> Result<()> {
let srv = start_500_server().await;
let coll = write_k6_script(&format!("http://{srv}"), "err500");
let config = JobConfig {
input: coll.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::ConstantVus {
vus: 2,
duration: "2s".to_string(),
graceful_stop: Some("3s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let result = engine.run(&config).await?;
let m = &result.metrics;
assert!(m.http_reqs > 0, "requests fired against the 500 server");
assert_eq!(
m.http_req_failed, 1.0,
"every request got a non-expected 500 -> http_req_failed == 1.0, got {}",
m.http_req_failed
);
let _ = std::fs::remove_file(&coll);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn hung_server_is_bounded_by_request_timeout() -> Result<()> {
let srv = start_hung_server().await;
let coll = write_k6_timeout_script(&format!("http://{srv}"), "hung", "500ms");
let config = JobConfig {
input: coll.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::ConstantVus {
vus: 2,
duration: "2s".to_string(),
graceful_stop: Some("3s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let start = Instant::now();
let result = engine.run(&config).await?;
let elapsed = start.elapsed();
let m = &result.metrics;
assert!(
elapsed < std::time::Duration::from_secs(20),
"run terminated despite hung server, took {elapsed:?}"
);
assert!(m.http_reqs > 0, "requests fired against the hung server");
assert_eq!(
m.http_req_failed, 1.0,
"every request timed out -> http_req_failed == 1.0, got {}",
m.http_req_failed
);
let _ = std::fs::remove_file(&coll);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn hung_server_is_bounded_by_global_request_timeout() -> Result<()> {
let srv = start_hung_server().await;
let coll = write_k6_script(&format!("http://{srv}"), "hungglobal");
let config = JobConfig {
input: coll.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::ConstantVus {
vus: 2,
duration: "2s".to_string(),
graceful_stop: Some("3s".to_string()),
think_time: ThinkTimeConfig::default(),
},
http: HttpConfig {
request_timeout: Some("500ms".to_string()),
..Default::default()
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let start = Instant::now();
let result = engine.run(&config).await?;
let elapsed = start.elapsed();
let m = &result.metrics;
assert!(
elapsed < std::time::Duration::from_secs(20),
"global request_timeout bounded the run, took {elapsed:?}"
);
assert!(m.http_reqs > 0, "requests fired against the hung server");
assert_eq!(
m.http_req_failed, 1.0,
"every request timed out -> http_req_failed == 1.0, got {}",
m.http_req_failed
);
let _ = std::fs::remove_file(&coll);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn connection_refused_is_recorded_as_failure() -> Result<()> {
let refused_addr = {
let listener = TcpListener::bind("127.0.0.1:0").await?;
let addr = listener.local_addr()?;
drop(listener);
addr
};
let coll = write_k6_script(&format!("http://{refused_addr}"), "refused");
let config = JobConfig {
input: coll.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::ConstantVus {
vus: 2,
duration: "2s".to_string(),
graceful_stop: Some("3s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let result = engine.run(&config).await?;
let m = &result.metrics;
assert_eq!(
m.http_req_failed, 1.0,
"connection refused -> http_req_failed == 1.0, got {}",
m.http_req_failed
);
assert!(m.http_reqs > 0, "requests fired against the refused port");
let _ = std::fs::remove_file(&coll);
Ok(())
}
fn write_k6_sleep_script(base: &str, tag: &str, sleep_secs: u64) -> String {
let dir = std::env::temp_dir();
let path = dir.join(format!("tropel-k6-e2e-{}-{}.js", std::process::id(), tag));
let script = format!(
r#"import http from 'k6/http';
import {{ check, sleep }} from 'k6';
export default function () {{
const res = http.get('{base}/');
check(res, {{ 'status is 200': (r) => r.status === 200 }});
sleep({sleep_secs});
}}
"#
);
std::fs::write(&path, script).unwrap();
path.to_string_lossy().to_string()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn force_stop_interrupts_sleeping_vu() -> Result<()> {
let srv = start_echo_server().await;
let coll = write_k6_sleep_script(&format!("http://{srv}"), "fsleep", 60);
let config = JobConfig {
input: coll.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::ConstantVus {
vus: 1,
duration: "2s".to_string(),
graceful_stop: Some("1s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let start = Instant::now();
let result = engine.run(&config).await?;
let elapsed = start.elapsed();
let m = &result.metrics;
assert!(
elapsed < std::time::Duration::from_secs(10),
"force-stop must interrupt the sleeping VU; run took {elapsed:?}"
);
assert!(m.http_reqs > 0, "requests fired, got {}", m.http_reqs);
let _ = std::fs::remove_file(&coll);
Ok(())
}
type SeenRequest = (String, Option<String>);
async fn start_cookie_server() -> (
std::net::SocketAddr,
Arc<std::sync::Mutex<Vec<SeenRequest>>>,
) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let seen: Arc<std::sync::Mutex<Vec<SeenRequest>>> = Arc::new(std::sync::Mutex::new(Vec::new()));
let seen_out = seen.clone();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
let seen = seen.clone();
tokio::spawn(async move {
let mut buf = [0u8; 4096];
loop {
let mut head = Vec::new();
loop {
let n = match sock.read(&mut buf).await {
Ok(0) | Err(_) => return,
Ok(n) => n,
};
head.extend_from_slice(&buf[..n]);
if head.windows(4).any(|w| w == b"\r\n\r\n") {
break;
}
}
let text = String::from_utf8_lossy(&head).into_owned();
let path = text
.lines()
.next()
.and_then(|l| l.split_whitespace().nth(1))
.unwrap_or("")
.to_string();
let cookie = text
.lines()
.find(|l| l.to_ascii_lowercase().starts_with("cookie:"))
.map(|l| l[7..].trim().to_string());
seen.lock().unwrap().push((path.clone(), cookie));
let body = r#"{"ok":true}"#;
let set = if path.starts_with("/set-cookie") {
"Set-Cookie: srv=from_server; Path=/\r\n"
} else {
""
};
let resp = format!(
"HTTP/1.1 200 OK\r\n{set}Content-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
);
if sock.write_all(resp.as_bytes()).await.is_err() {
return;
}
}
});
}
});
(addr, seen_out)
}
fn write_k6_cookie_jar_script(base: &str, tag: &str) -> String {
let dir = std::env::temp_dir();
let path = dir.join(format!("tropel-k6-e2e-{}-{}.js", std::process::id(), tag));
let script = format!(
r#"import http from 'k6/http';
export default function () {{
const jar = http.cookieJar();
// (1) A cookie the SCRIPT set must ride on the next request.
jar.set('{base}/', 'scripted', 'from_jar', {{ path: '/' }});
http.get('{base}/first');
// (2) The server sets one of its own.
http.get('{base}/set-cookie');
// (3) …which the script must be able to READ back out of the same jar.
// Whatever it saw becomes the path, so the server records it.
const seen = jar.cookiesForURL('{base}/');
const srv = (seen && seen['srv'] && seen['srv'][0]) ? seen['srv'][0] : 'MISSING';
http.get('{base}/probe-' + srv);
// (4) delete() drops ONLY the named cookie — the server's must survive.
jar.delete('{base}/', 'scripted');
http.get('{base}/after-delete');
}}
"#
);
std::fs::write(&path, script).unwrap();
path.to_string_lossy().to_string()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn cookie_jar_set_reaches_the_wire_and_reads_back_server_cookies() -> Result<()> {
let (srv, seen) = start_cookie_server().await;
let script = write_k6_cookie_jar_script(&format!("http://{srv}"), "cookiejar");
let config = JobConfig {
input: script.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::SharedIterations {
iterations: 1,
max_duration: Some("30s".to_string()),
vus: 1,
graceful_stop: Some("5s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
let result = engine.run(&config).await?;
assert!(
result.metrics.http_reqs >= 4,
"the script's four requests must all fire, got {}",
result.metrics.http_reqs
);
let seen = seen.lock().unwrap().clone();
let find = |p: &str| -> SeenRequest {
seen.iter()
.find(|(path, _)| path == p)
.unwrap_or_else(|| {
panic!(
"the server never received {p}; it saw: {:?}",
seen.iter().map(|(p, _)| p.as_str()).collect::<Vec<_>>()
)
})
.clone()
};
let (_, first_cookie) = find("/first");
let first_cookie = first_cookie.unwrap_or_else(|| {
panic!("jar.set() must put a Cookie header on the next request; none was sent")
});
assert!(
first_cookie.contains("scripted=from_jar"),
"jar.set() must reach the wire; server saw Cookie: {first_cookie}"
);
let (probe_path, _) = seen
.iter()
.find(|(p, _)| p.starts_with("/probe-"))
.unwrap_or_else(|| {
panic!(
"the server never received a /probe- request; it saw: {:?}",
seen.iter().map(|(p, _)| p.as_str()).collect::<Vec<_>>()
)
})
.clone();
assert_eq!(
probe_path, "/probe-from_server",
"cookiesForURL() must return the server's cookie VALUE keyed by name \
(k6 returns map[string][]string, so cookies['srv'][0]); \
`/probe-MISSING` means the jar read nothing back"
);
let (_, after_delete) = find("/after-delete");
let after_delete = after_delete.unwrap_or_else(|| {
panic!("the server's own cookie must survive delete('scripted'); no Cookie header was sent")
});
assert!(
!after_delete.contains("scripted="),
"delete() must stop the named cookie being sent; server saw Cookie: {after_delete}"
);
assert!(
after_delete.contains("srv=from_server"),
"delete() must NOT drop other cookies; server saw Cookie: {after_delete}"
);
let _ = std::fs::remove_file(&script);
Ok(())
}
async fn start_inflight_peak_server(
hold_ms: u64,
) -> (std::net::SocketAddr, Arc<AtomicUsize>, Arc<AtomicUsize>) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let active = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let total = Arc::new(AtomicUsize::new(0));
let (peak_out, total_out) = (peak.clone(), total.clone());
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
let (active, peak, total) = (active.clone(), peak.clone(), total.clone());
tokio::spawn(async move {
let mut buf = [0u8; 4096];
loop {
let mut head = Vec::new();
loop {
let n = match sock.read(&mut buf).await {
Ok(0) | Err(_) => return,
Ok(n) => n,
};
head.extend_from_slice(&buf[..n]);
if head.windows(4).any(|w| w == b"\r\n\r\n") {
break;
}
}
total.fetch_add(1, Ordering::SeqCst);
let now = active.fetch_add(1, Ordering::SeqCst) + 1;
peak.fetch_max(now, Ordering::SeqCst);
tokio::time::sleep(std::time::Duration::from_millis(hold_ms)).await;
let body = r#"{"ok":true}"#;
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
);
let wrote = sock.write_all(resp.as_bytes()).await;
active.fetch_sub(1, Ordering::SeqCst);
if wrote.is_err() {
return;
}
}
});
}
});
(addr, peak_out, total_out)
}
fn write_k6_batch_script(base: &str, tag: &str, batch: usize, per_host: usize, n: usize) -> String {
let dir = std::env::temp_dir();
let path = dir.join(format!("tropel-k6-e2e-{}-{}.js", std::process::id(), tag));
let script = format!(
r#"import http from 'k6/http';
export const options = {{ batch: {batch}, batchPerHost: {per_host} }};
export default function () {{
const reqs = [];
for (let i = 0; i < {n}; i++) {{ reqs.push(['GET', '{base}/b' + i]); }}
http.batch(reqs);
}}
"#
);
std::fs::write(&path, script).unwrap();
path.to_string_lossy().to_string()
}
async fn observed_batch_concurrency(
tag: &str,
batch: usize,
per_host: usize,
n: usize,
) -> Result<(usize, usize)> {
let (srv, peak, total) = start_inflight_peak_server(150).await;
let script = write_k6_batch_script(&format!("http://{srv}"), tag, batch, per_host, n);
let config = JobConfig {
input: script.clone(),
input_type: Some("k6".to_string()),
execution: ExecutionConfig::SharedIterations {
iterations: 1,
max_duration: Some("60s".to_string()),
vus: 1,
graceful_stop: Some("10s".to_string()),
think_time: ThinkTimeConfig::default(),
},
output: OutputConfig {
reporters: vec![],
..Default::default()
},
..Default::default()
};
let engine = Engine::new(ExtensionRegistry::new());
engine.run(&config).await?;
let _ = std::fs::remove_file(&script);
Ok((peak.load(Ordering::SeqCst), total.load(Ordering::SeqCst)))
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn declared_batch_limit_is_the_observed_concurrency() -> Result<()> {
let (peak_low, total_low) = observed_batch_concurrency("batch-low", 2, 2, 12).await?;
assert_eq!(total_low, 12, "all 12 batch requests must be served");
assert_eq!(
peak_low, 2,
"declared `batch: 2` must be the observed ceiling; the server saw \
{peak_low} requests in flight at once (6 = k6's per-host default \
still hardcoded, i.e. the declared option was dropped)"
);
let (peak_high, total_high) = observed_batch_concurrency("batch-high", 10, 10, 12).await?;
assert_eq!(total_high, 12, "all 12 batch requests must be served");
assert_eq!(
peak_high, 10,
"declared `batch: 10` must raise the ceiling above k6's per-host \
default of 6; the server saw {peak_high} in flight at once"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn batch_per_host_caps_below_the_global_batch_limit() -> Result<()> {
let (peak, total) = observed_batch_concurrency("batch-perhost", 12, 3, 12).await?;
assert_eq!(total, 12, "all 12 batch requests must be served");
assert_eq!(
peak, 3,
"`batchPerHost: 3` must cap a single-host batch below `batch: 12`; \
the server saw {peak} in flight at once"
);
Ok(())
}