#![cfg(test)]
mod common;
use common::soak::{parse_duration_from_args, record_latency, SoakMonitor, SoakRegression};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use sz_rust_core::cache::{Cache, MemoryCacheDriver};
use sz_rust_core::container::Container;
use sz_rust_core::response::ApiResponse;
use sz_rust_core::router::parse_path;
use sz_rust_core::routing::HandlerRef;
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "soak test 需显式 --ignored 启动;默认 60s,6h 任务 SOAK_DURATION=6h"]
async fn soak_web_framework_steady_state() {
let duration = parse_duration_from_args();
let sample_interval = Duration::from_secs(if duration.as_secs() >= 3600 {
60 } else {
5 });
eprintln!(
"[soak] 启动:duration={:?}, sample_interval={:?}",
duration, sample_interval
);
let mut monitor = SoakMonitor::new(duration, sample_interval);
let ops_counter = monitor.ops_counter();
let errors_counter = monitor.errors_counter();
let latency_window = monitor.latency_window();
const WORKER_COUNT: usize = 8;
let stop_flag = Arc::new(AtomicBool::new(false));
let mut workers = Vec::new();
for worker_id in 0..WORKER_COUNT {
let ops_clone = ops_counter.clone();
let errors_clone = errors_counter.clone();
let latency_clone = latency_window.clone();
let stop_clone = stop_flag.clone();
workers.push(tokio::spawn(async move {
let uris = [
"/oapc/customer/index",
"/admin/login/index",
"/api/user/list",
"/oapc/order/detail",
"/farm/animal/feed",
];
let handlers = [
"Customer@index",
"Login@index",
"User@list",
"Order@detail",
"Animal@feed",
];
let chain = sz_rust_core::middleware::chain::MiddlewareChain::default_chain();
let container = Container::new();
container.singleton(|| 42i32);
let cache = Arc::new(Cache::new());
cache.register_default(MemoryCacheDriver::new());
let cache_key = format!("soak_worker_{}", worker_id);
let mut iteration: usize = 0;
while !stop_clone.load(Ordering::Relaxed) {
let t0 = Instant::now();
let uri = uris[iteration % uris.len()];
let parsed = parse_path(uri);
let handler_str = handlers[iteration % handlers.len()];
let handler = match HandlerRef::parse(handler_str) {
Ok(h) => h,
Err(e) => {
errors_clone.fetch_add(1, Ordering::Relaxed);
eprintln!(
"[soak worker {}] HandlerRef::parse({}) error: {}",
worker_id, handler_str, e
);
tokio::time::sleep(Duration::from_millis(10)).await;
continue;
}
};
let resp = ApiResponse::success(
serde_json::json!({
"app": parsed.app,
"controller": parsed.controller,
"action": parsed.action,
"handler": handler.to_handler_string(),
}),
"ok",
);
let _json = resp.to_json_string();
let _has_dup = chain.has_duplicates();
let _val = container.make::<i32>();
let _ = cache.set(&cache_key, iteration, None);
let _: Option<usize> = cache.get(&cache_key).unwrap_or(None);
let elapsed_us = t0.elapsed().as_micros() as u64;
record_latency(&latency_clone, elapsed_us);
ops_clone.fetch_add(1, Ordering::Relaxed);
iteration = iteration.wrapping_add(1);
tokio::task::yield_now().await;
}
}));
}
while !monitor.is_finished() {
tokio::time::sleep(sample_interval).await;
let snap = monitor.snapshot((0, 0, 0));
eprintln!(
"[soak] t={}s ops={} ops/s={:.1} rss={}MB fd={} threads={} p99={}us errors={}",
snap.elapsed_secs,
snap.ops_completed,
snap.ops_per_sec,
snap.rss_bytes / 1024 / 1024,
snap.fd_count,
snap.thread_count,
snap.p99_latency_us,
snap.error_count,
);
}
stop_flag.store(true, Ordering::Release);
for w in workers {
let _ = w.await;
}
let final_snap = monitor.snapshot((0, 0, 0));
eprintln!(
"[soak] 完成:总操作 {} 次,错误 {} 次",
final_snap.ops_completed, final_snap.error_count,
);
let csv_path = std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("..")
.join("..")
.join("target")
.join("soak-report.csv");
let csv_path_str = csv_path.to_str().expect("CSV path not UTF-8");
if let Err(e) = monitor.export_csv(csv_path_str) {
eprintln!("[soak] CSV 导出失败: {}", e);
} else {
eprintln!("[soak] CSV 报告已导出: {}", csv_path_str);
}
let regressions = monitor.detect_regressions();
if regressions.is_empty() {
eprintln!("[soak] ✅ 未检测到退化");
} else {
eprintln!("[soak] ⚠ 检测到 {} 项退化:", regressions.len());
for r in ®ressions {
eprintln!(" - {}", r);
}
}
assert!(
regressions.is_empty(),
"Soak test 检测到退化:{}",
regressions
.iter()
.map(|r| r.to_string())
.collect::<Vec<_>>()
.join("; ")
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn soak_smoke_10s() {
let duration = Duration::from_secs(10);
let sample_interval = Duration::from_secs(2);
let mut monitor = SoakMonitor::new(duration, sample_interval);
let ops = monitor.ops_counter();
let errors = monitor.errors_counter();
let latency = monitor.latency_window();
let stop = Arc::new(AtomicBool::new(false));
let mut workers = Vec::new();
for _ in 0..4 {
let ops_c = ops.clone();
let err_c = errors.clone();
let lat_c = latency.clone();
let stop_c = stop.clone();
workers.push(tokio::spawn(async move {
let uris = [
"/oapc/customer/index",
"/admin/login/index",
"/api/user/list",
];
let handlers = ["Customer@index", "Login@index", "User@list"];
let chain = sz_rust_core::middleware::chain::MiddlewareChain::default_chain();
let container = Container::new();
container.singleton(|| 42i32);
let cache = Arc::new(Cache::new());
cache.register_default(MemoryCacheDriver::new());
let mut iteration: usize = 0;
while !stop_c.load(Ordering::Relaxed) {
let t0 = Instant::now();
let uri = uris[iteration % uris.len()];
let parsed = parse_path(uri);
let handler_str = handlers[iteration % handlers.len()];
if let Ok(handler) = HandlerRef::parse(handler_str) {
let resp = ApiResponse::success(
serde_json::json!({
"app": parsed.app,
"controller": parsed.controller,
"action": parsed.action,
"handler": handler.to_handler_string(),
}),
"ok",
);
let _json = resp.to_json_string();
let _has_dup = chain.has_duplicates();
let _val = container.make::<i32>();
let _ = cache.set("soak_smoke", iteration, None);
let _: Option<usize> = cache.get("soak_smoke").unwrap_or(None);
record_latency(&lat_c, t0.elapsed().as_micros() as u64);
ops_c.fetch_add(1, Ordering::Relaxed);
} else {
err_c.fetch_add(1, Ordering::Relaxed);
}
iteration = iteration.wrapping_add(1);
tokio::task::yield_now().await;
}
}));
}
while !monitor.is_finished() {
tokio::time::sleep(sample_interval).await;
let snap = monitor.snapshot((0, 0, 0));
eprintln!(
"[soak-smoke] t={}s ops={} rss={}MB p99={}us errors={}",
snap.elapsed_secs,
snap.ops_completed,
snap.rss_bytes / 1024 / 1024,
snap.p99_latency_us,
snap.error_count,
);
}
stop.store(true, Ordering::Release);
for w in workers {
let _ = w.await;
}
monitor.snapshot((0, 0, 0));
let total_ops = monitor
.snapshots()
.last()
.map(|s| s.ops_completed)
.unwrap_or(0);
assert!(
total_ops >= 50,
"Soak smoke 应至少完成 50 次操作,实际 {}",
total_ops
);
let regressions = monitor.detect_regressions();
let critical: Vec<&SoakRegression> = regressions
.iter()
.filter(|r| matches!(r, SoakRegression::PoolLeak { .. }))
.collect();
assert!(
critical.is_empty(),
"Soak smoke 不应有 PoolLeak:{:?}",
critical
);
}