use std::sync::{
Arc, LazyLock,
atomic::{AtomicUsize, Ordering},
};
use tokio::{
io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader},
net::TcpListener,
};
static BRIDGE: LazyLock<agy_bridge::AgyBridge> = LazyLock::new(|| {
agy_bridge::AgyBridge::builder()
.inter_agent_delay(std::time::Duration::ZERO)
.build()
.expect("shared AgyBridge")
});
async fn parse_http_request<R: tokio::io::AsyncRead + Unpin>(
buf_reader: &mut BufReader<R>,
) -> Option<String> {
let mut request_line = String::new();
if let Err(e) = buf_reader.read_line(&mut request_line).await {
eprintln!("mock server: failed to read request line: {e}");
return None;
}
let request_line = request_line.trim_end().to_string();
if request_line.is_empty() {
return None;
}
let mut content_length: usize = 0;
loop {
let mut line = String::new();
if let Err(e) = buf_reader.read_line(&mut line).await {
eprintln!("mock server: failed to read header line: {e}");
return None;
}
let trimmed = line.trim();
if trimmed.is_empty() {
break;
}
let lower = trimmed.to_lowercase();
if let Some(val) = lower.strip_prefix("content-length:") {
content_length = match val.trim().parse() {
Ok(len) => len,
Err(e) => {
eprintln!("mock server: invalid Content-Length header: {e}");
return None;
}
};
}
}
if content_length > 0 {
let mut body_buf = vec![0u8; content_length];
if let Err(e) = buf_reader.read_exact(&mut body_buf).await {
eprintln!("mock server: failed to read body: {e}");
return None;
}
}
Some(request_line)
}
fn json_response(status: u16, body: &str) -> String {
let reason = match status {
200 => "OK",
404 => "Not Found",
429 => "Too Many Requests",
500 => "Internal Server Error",
503 => "Service Unavailable",
_ => "Error",
};
format!(
"HTTP/1.1 {status} {reason}\r\n\
Content-Type: application/json\r\n\
Content-Length: {}\r\n\
\r\n\
{}",
body.len(),
body
)
}
fn sse_response(json_body: &str) -> String {
let sse_data = format!("data: {json_body}\n\n");
format!(
"HTTP/1.1 200 OK\r\n\
Content-Type: text/event-stream\r\n\
Content-Length: {}\r\n\
\r\n\
{}",
sse_data.len(),
sse_data
)
}
fn model_list_json() -> String {
serde_json::json!({
"models": [
{
"name": "models/gemini-3.6-flash",
"displayName": "Gemini 3.6 Flash",
"supportedGenerationMethods": [
"generateContent",
"streamGenerateContent",
"countTokens"
],
"inputTokenLimit": 1_048_576,
"outputTokenLimit": 8192
},
{
"name": "models/gemini-3.5-flash",
"displayName": "Gemini 3.5 Flash",
"supportedGenerationMethods": [
"generateContent",
"streamGenerateContent",
"countTokens"
],
"inputTokenLimit": 1_048_576,
"outputTokenLimit": 8192
},
{
"name": "models/gemini-2.0-flash",
"displayName": "Gemini 2.0 Flash",
"supportedGenerationMethods": [
"generateContent",
"streamGenerateContent",
"countTokens"
],
"inputTokenLimit": 1_048_576,
"outputTokenLimit": 8192
}
]
})
.to_string()
}
fn generate_content_json(tag: &str) -> String {
serde_json::json!({
"candidates": [{
"content": {
"parts": [{"text": format!("mock:{tag}")}],
"role": "model"
},
"finishReason": "STOP",
"index": 0
}],
"usageMetadata": {
"promptTokenCount": 10,
"candidatesTokenCount": 5,
"totalTokenCount": 15
}
})
.to_string()
}
struct MockServer {
addr: std::net::SocketAddr,
post_count: Arc<AtomicUsize>,
handle: tokio::task::JoinHandle<()>,
}
async fn serve_connection(
stream: tokio::net::TcpStream,
tag: String,
count: Arc<AtomicUsize>,
delay: Option<std::time::Duration>,
) {
let (reader, mut writer) = tokio::io::split(stream);
let mut buf_reader = BufReader::new(reader);
loop {
let Some(request_line) = parse_http_request(&mut buf_reader).await else {
break;
};
let response = if request_line.starts_with("GET ") {
json_response(200, &model_list_json())
} else {
count.fetch_add(1, Ordering::SeqCst);
if let Some(delay) = delay {
tokio::time::sleep(delay).await;
}
sse_response(&generate_content_json(&tag))
};
if let Err(e) = writer.write_all(response.as_bytes()).await {
eprintln!("[MOCK {tag}] write error: {e}");
break;
}
if let Err(e) = writer.flush().await {
eprintln!("[MOCK {tag}] flush error: {e}");
break;
}
}
}
impl MockServer {
async fn start(tag: &str) -> Self {
Self::start_with_delay(tag, None).await
}
async fn start_slow(tag: &str, delay: std::time::Duration) -> Self {
Self::start_with_delay(tag, Some(delay)).await
}
async fn start_with_delay(tag: &str, delay: Option<std::time::Duration>) -> Self {
let listener = TcpListener::bind("127.0.0.1:0")
.await
.expect("bind mock server");
let addr = listener.local_addr().expect("local addr");
let post_count = Arc::new(AtomicUsize::new(0));
let count = Arc::clone(&post_count);
let tag = tag.to_string();
let handle = tokio::spawn(async move {
loop {
let Ok((stream, _)) = listener.accept().await else {
break;
};
tokio::spawn(serve_connection(
stream,
tag.clone(),
Arc::clone(&count),
delay,
));
}
});
Self {
addr,
post_count,
handle,
}
}
fn base_url(&self) -> String {
format!("http://{}", self.addr)
}
fn post_count(&self) -> usize {
self.post_count.load(Ordering::SeqCst)
}
}
impl Drop for MockServer {
fn drop(&mut self) {
self.handle.abort();
}
}
fn agent_config(base_url: &str, system: &str) -> agy_bridge::config::AgentConfig {
agy_bridge::config::AgentConfig::builder()
.system_instructions(system)
.gemini(agy_bridge::config::GeminiConfig {
api_key: Some("test-key".to_string()),
base_url: Some(base_url.to_string()),
models: agy_bridge::config::ModelConfig::default(),
})
.capabilities(agy_bridge::config::CapabilitiesConfig::custom_tools_only())
.retry_config(agy_bridge::config::RetryConfig::no_retries())
.build()
}
fn multi_thread_rt() -> tokio::runtime::Runtime {
tokio::runtime::Builder::new_multi_thread()
.worker_threads(4)
.enable_all()
.build()
.expect("multi-thread tokio runtime")
}
async fn wait_for_zero_agents(bridge: &agy_bridge::AgyBridge) -> bool {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
loop {
match bridge.active_agent_count().await {
Ok(0) => return true,
Ok(_) => {}
Err(e) => panic!("active_agent_count query failed: {e}"),
}
if std::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
}
async fn warm_up_python_sdk(base_url: &str) {
let warmup = BRIDGE
.agent(agent_config(base_url, "warmup"))
.await
.expect("warm-up agent creation");
warmup.shutdown().await.expect("warm-up agent shutdown");
}
#[test]
fn concurrent_agent_startup_no_gil_lockup() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockServer::start("startup").await;
let url = server.base_url();
warm_up_python_sdk(&url).await;
let start = std::time::Instant::now();
let (a1, a2, a3, a4, a5) =
tokio::time::timeout(std::time::Duration::from_secs(30), async {
tokio::join!(
BRIDGE.agent(agent_config(&url, "Agent 1")),
BRIDGE.agent(agent_config(&url, "Agent 2")),
BRIDGE.agent(agent_config(&url, "Agent 3")),
BRIDGE.agent(agent_config(&url, "Agent 4")),
BRIDGE.agent(agent_config(&url, "Agent 5")),
)
})
.await
.expect("concurrent agent creation stalled — possible GIL deadlock");
let elapsed = start.elapsed();
eprintln!("5 concurrent agent creations took {elapsed:.1?}");
let a1 = a1.expect("agent 1");
let a2 = a2.expect("agent 2");
let a3 = a3.expect("agent 3");
let a4 = a4.expect("agent 4");
let a5 = a5.expect("agent 5");
assert!(
elapsed.as_secs() < 14,
"5 concurrent creates took {elapsed:.1?} — possible GIL deadlock"
);
let (r1, r2, r3, r4, r5) = tokio::join!(
a1.chat_text("ping"),
a2.chat_text("ping"),
a3.chat_text("ping"),
a4.chat_text("ping"),
a5.chat_text("ping"),
);
assert!(r1.is_ok(), "Agent 1 chat failed: {r1:?}");
assert!(r2.is_ok(), "Agent 2 chat failed: {r2:?}");
assert!(r3.is_ok(), "Agent 3 chat failed: {r3:?}");
assert!(r4.is_ok(), "Agent 4 chat failed: {r4:?}");
assert!(r5.is_ok(), "Agent 5 chat failed: {r5:?}");
a1.shutdown().await.expect("shutdown a1");
a2.shutdown().await.expect("shutdown a2");
a3.shutdown().await.expect("shutdown a3");
a4.shutdown().await.expect("shutdown a4");
a5.shutdown().await.expect("shutdown a5");
});
}
#[test]
fn ongoing_chat_survives_peer_shutdown() {
let rt = multi_thread_rt();
rt.block_on(async {
let fast_server = MockServer::start("fast").await;
let slow_server =
MockServer::start_slow("slow", std::time::Duration::from_millis(500)).await;
let agent_a = BRIDGE
.agent(agent_config(&fast_server.base_url(), "fast-agent"))
.await
.expect("agent A");
let agent_b = BRIDGE
.agent(agent_config(&slow_server.base_url(), "slow-agent"))
.await
.expect("agent B");
let warmup = agent_a.chat_text("hello").await;
assert!(warmup.is_ok(), "warmup failed: {warmup:?}");
let (b_result, a_shutdown) =
tokio::join!(agent_b.chat_text("slow request"), agent_a.shutdown(),);
a_shutdown.expect("agent A shutdown");
let b_text = b_result.expect("agent B chat should succeed during A's shutdown");
assert!(
b_text.contains("mock:slow"),
"Expected slow mock response, got: {b_text}"
);
agent_b.shutdown().await.expect("agent B shutdown");
});
}
#[test]
fn new_agent_during_peer_shutdown() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockServer::start("create-during-shutdown").await;
let url = server.base_url();
let agent_a = BRIDGE
.agent(agent_config(&url, "A"))
.await
.expect("agent A");
agent_a.chat_text("warmup").await.expect("A warmup");
let (a_shutdown, c_creation) = tokio::join!(
agent_a.shutdown(),
BRIDGE.agent(agent_config(&url, "C-new")),
);
a_shutdown.expect("A shutdown");
let agent_c = c_creation.expect("agent C creation during A shutdown");
let c_text = agent_c.chat_text("hello from C").await.expect("C chat");
assert!(
c_text.contains("mock:create-during-shutdown"),
"Expected mock response from C, got: {c_text}"
);
agent_c.shutdown().await.expect("C shutdown");
});
}
#[test]
fn agents_work_after_full_teardown() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockServer::start("post-teardown").await;
let url = server.base_url();
{
let agent = BRIDGE
.agent(agent_config(&url, "phase1"))
.await
.expect("phase1 agent");
let text = agent.chat_text("phase1").await.expect("phase1 chat");
assert!(text.contains("mock:post-teardown"), "phase1 got: {text}");
agent.shutdown().await.expect("phase1 shutdown");
}
let agent_new = BRIDGE
.agent(agent_config(&url, "phase2"))
.await
.expect("phase2 agent");
let text = agent_new.chat_text("phase2").await.expect("phase2 chat");
assert!(text.contains("mock:post-teardown"), "phase2 got: {text}");
agent_new.shutdown().await.expect("phase2 shutdown");
assert!(
server.post_count() >= 2,
"Expected at least 2 POST requests (one per phase), got {}",
server.post_count()
);
});
}
#[test]
fn no_python_object_leaks_rapid_cycles() {
let rt = multi_thread_rt();
rt.block_on(async {
const CYCLES: usize = 10;
let bridge = agy_bridge::AgyBridge::builder()
.inter_agent_delay(std::time::Duration::ZERO)
.build()
.expect("dedicated leak-test bridge");
let server = MockServer::start("leak-test").await;
let url = server.base_url();
assert_eq!(
bridge
.active_agent_count()
.await
.expect("count on fresh bridge"),
0,
"a fresh bridge must have no live agents"
);
for i in 0..CYCLES {
let agent = bridge
.agent(agent_config(&url, &format!("cycle-{i}")))
.await
.unwrap_or_else(|e| panic!("agent creation failed at cycle {i}: {e}"));
assert_eq!(
bridge
.active_agent_count()
.await
.expect("count with live agent"),
1,
"exactly one agent should be live during cycle {i}"
);
let text = agent
.chat_text(format!("cycle {i}"))
.await
.unwrap_or_else(|e| panic!("chat failed at cycle {i}: {e}"));
assert!(text.contains("mock:leak-test"), "cycle {i} got: {text}");
agent
.shutdown()
.await
.unwrap_or_else(|e| panic!("shutdown failed at cycle {i}: {e}"));
let live = bridge
.active_agent_count()
.await
.expect("count after shutdown");
assert_eq!(
live, 0,
"agent leak after cycle {i}: {live} agent(s) still registered"
);
}
{
let agent = bridge
.agent(agent_config(&url, "drop-cycle"))
.await
.expect("drop-cycle agent");
assert_eq!(
bridge
.active_agent_count()
.await
.expect("count before drop"),
1,
"one live agent before drop"
);
drop(agent);
}
assert!(
wait_for_zero_agents(&bridge).await,
"agent dropped without shutdown() was not cleaned up — leak"
);
assert_eq!(
server.post_count(),
CYCLES,
"expected exactly {CYCLES} POST requests"
);
});
}
#[test]
fn sequential_teardown_while_others_chat() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockServer::start("seq-teardown").await;
let url = server.base_url();
let a = BRIDGE.agent(agent_config(&url, "A")).await.expect("A");
let b = BRIDGE.agent(agent_config(&url, "B")).await.expect("B");
let c = BRIDGE.agent(agent_config(&url, "C")).await.expect("C");
let (ra, rb, rc) = tokio::join!(
a.chat_text("hello A"),
b.chat_text("hello B"),
c.chat_text("hello C"),
);
assert!(ra.is_ok(), "A chat failed: {ra:?}");
assert!(rb.is_ok(), "B chat failed: {rb:?}");
assert!(rc.is_ok(), "C chat failed: {rc:?}");
a.shutdown().await.expect("A shutdown");
let rb2 = b.chat_text("after A shutdown").await;
let rc2 = c.chat_text("after A shutdown").await;
assert!(rb2.is_ok(), "B should work after A shutdown: {rb2:?}");
assert!(rc2.is_ok(), "C should work after A shutdown: {rc2:?}");
b.shutdown().await.expect("B shutdown");
let rc3 = c.chat_text("after B shutdown").await;
assert!(rc3.is_ok(), "C should work after B shutdown: {rc3:?}");
c.shutdown().await.expect("C shutdown");
});
}
#[test]
fn two_bridges_concurrent_agents_isolation() {
let rt = multi_thread_rt();
rt.block_on(async {
let server_1 = MockServer::start("bridge1").await;
let server_2 = MockServer::start("bridge2").await;
let bridge_1 = agy_bridge::AgyBridge::builder()
.inter_agent_delay(std::time::Duration::ZERO)
.build()
.expect("bridge 1");
let bridge_2 = agy_bridge::AgyBridge::builder()
.inter_agent_delay(std::time::Duration::ZERO)
.build()
.expect("bridge 2");
let a1 = bridge_1
.agent(agent_config(&server_1.base_url(), "b1-agent"))
.await
.expect("b1 agent");
let a2 = bridge_2
.agent(agent_config(&server_2.base_url(), "b2-agent"))
.await
.expect("b2 agent");
let (r1, r2) = tokio::join!(a1.chat_text("b1 ping"), a2.chat_text("b2 ping"),);
let t1 = r1.expect("bridge 1 chat");
let t2 = r2.expect("bridge 2 chat");
assert!(t1.contains("mock:bridge1"), "Bridge 1 response wrong: {t1}");
assert!(t2.contains("mock:bridge2"), "Bridge 2 response wrong: {t2}");
assert_eq!(server_1.post_count(), 1, "bridge 1 should get 1 POST");
assert_eq!(server_2.post_count(), 1, "bridge 2 should get 1 POST");
a1.shutdown().await.expect("b1 shutdown");
let r2_after = a2.chat_text("after b1 shutdown").await;
r2_after.expect("Bridge 2 agent should survive bridge 1 teardown");
a2.shutdown().await.expect("b2 shutdown");
});
}
#[test]
fn concurrent_creation_across_bridges_no_interference() {
let rt = multi_thread_rt();
rt.block_on(async {
const PER_BRIDGE: usize = 4;
let server = MockServer::start("multi-bridge").await;
let url = server.base_url();
let bridge_1 = agy_bridge::AgyBridge::builder()
.inter_agent_delay(std::time::Duration::ZERO)
.build()
.expect("bridge 1");
let bridge_2 = agy_bridge::AgyBridge::builder()
.inter_agent_delay(std::time::Duration::ZERO)
.build()
.expect("bridge 2");
let (b1, b2) = (&bridge_1, &bridge_2);
let (created_1, created_2) = tokio::join!(
futures::future::join_all((0..PER_BRIDGE).map(|i| {
let url = url.clone();
async move { b1.agent(agent_config(&url, &format!("b1-{i}"))).await }
})),
futures::future::join_all((0..PER_BRIDGE).map(|i| {
let url = url.clone();
async move { b2.agent(agent_config(&url, &format!("b2-{i}"))).await }
})),
);
let agents_1: Vec<_> = created_1
.into_iter()
.enumerate()
.map(|(i, r)| r.unwrap_or_else(|e| panic!("bridge 1 agent {i} creation: {e}")))
.collect();
let agents_2: Vec<_> = created_2
.into_iter()
.enumerate()
.map(|(i, r)| r.unwrap_or_else(|e| panic!("bridge 2 agent {i} creation: {e}")))
.collect();
assert_eq!(
bridge_1.active_agent_count().await.expect("b1 count"),
PER_BRIDGE,
"bridge 1 must see exactly its own agents"
);
assert_eq!(
bridge_2.active_agent_count().await.expect("b2 count"),
PER_BRIDGE,
"bridge 2 must see exactly its own agents"
);
let ids: Vec<u64> = agents_1
.iter()
.chain(agents_2.iter())
.map(agy_bridge::agent::AgentHandle::id)
.collect();
let unique: std::collections::HashSet<u64> = ids.iter().copied().collect();
assert_eq!(
ids.len(),
unique.len(),
"agent IDs must be globally unique across bridges: {ids:?}"
);
futures::future::join_all(
agents_1
.into_iter()
.map(|a| async move { a.shutdown().await }),
)
.await
.into_iter()
.enumerate()
.for_each(|(i, r)| r.unwrap_or_else(|e| panic!("bridge 1 shutdown {i}: {e}")));
assert_eq!(
bridge_1.active_agent_count().await.expect("b1 count after"),
0,
"bridge 1 must be drained after shutting down all its agents"
);
assert_eq!(
bridge_2.active_agent_count().await.expect("b2 count after"),
PER_BRIDGE,
"bridge 2 agents must survive bridge 1's full teardown"
);
for (i, a) in agents_2.iter().enumerate() {
let text = a
.chat_text(format!("b2 still alive {i}"))
.await
.unwrap_or_else(|e| panic!("bridge 2 chat {i} after bridge 1 teardown: {e}"));
assert!(
text.contains("mock:multi-bridge"),
"bridge 2 agent {i}: {text}"
);
}
futures::future::join_all(
agents_2
.into_iter()
.map(|a| async move { a.shutdown().await }),
)
.await
.into_iter()
.enumerate()
.for_each(|(i, r)| r.unwrap_or_else(|e| panic!("bridge 2 shutdown {i}: {e}")));
assert!(
wait_for_zero_agents(&bridge_2).await,
"bridge 2 must drain to zero agents after teardown"
);
});
}
#[test]
fn sequential_agent_startup_timing() {
let rt = multi_thread_rt();
rt.block_on(async {
const COUNT: usize = 5;
let server = MockServer::start("timing").await;
let url = server.base_url();
let mut agents = Vec::with_capacity(COUNT);
warm_up_python_sdk(&url).await;
for i in 0..COUNT {
let start = std::time::Instant::now();
let agent = BRIDGE
.agent(agent_config(&url, &format!("timing-{i}")))
.await
.unwrap_or_else(|e| panic!("agent {i} creation failed: {e}"));
let elapsed = start.elapsed();
eprintln!("Agent {i} creation: {elapsed:.1?}");
assert!(
elapsed.as_secs() < 14,
"Agent {i} creation took {elapsed:.1?} — possible GIL contention"
);
agents.push(agent);
}
for agent in agents {
agent.shutdown().await.expect("shutdown");
}
});
}
#[test]
fn create_agent_during_peer_active_chat() {
let rt = multi_thread_rt();
rt.block_on(async {
let slow_server =
MockServer::start_slow("slow-chat", std::time::Duration::from_millis(500)).await;
let fast_server = MockServer::start("fast-create").await;
let agent_a = BRIDGE
.agent(agent_config(&slow_server.base_url(), "slow-chatter"))
.await
.expect("agent A");
let (a_chat_result, b_creation_result) = tokio::join!(
agent_a.chat_text("slow request"),
BRIDGE.agent(agent_config(&fast_server.base_url(), "fast-new")),
);
let a_text = a_chat_result.expect("A chat during concurrent B creation");
assert!(a_text.contains("mock:slow-chat"), "A got: {a_text}");
let agent_b = b_creation_result.expect("B creation during A's active chat");
let b_text = agent_b.chat_text("hello B").await.expect("B chat");
assert!(b_text.contains("mock:fast-create"), "B got: {b_text}");
agent_a.shutdown().await.expect("A shutdown");
agent_b.shutdown().await.expect("B shutdown");
});
}
#[test]
fn rapid_concurrent_lifecycle_stress() {
let rt = multi_thread_rt();
rt.block_on(async {
const N: usize = 4;
let server = MockServer::start("stress").await;
let url = server.base_url();
let mut handles = Vec::with_capacity(N);
for i in 0..N {
let url = url.clone();
handles.push(tokio::spawn(async move {
let agent = BRIDGE
.agent(agent_config(&url, &format!("stress-{i}")))
.await
.unwrap_or_else(|e| panic!("stress agent {i} creation: {e}"));
let text = agent
.chat_text(format!("stress ping {i}"))
.await
.unwrap_or_else(|e| panic!("stress agent {i} chat: {e}"));
assert!(text.contains("mock:stress"), "stress agent {i} got: {text}");
agent
.shutdown()
.await
.unwrap_or_else(|e| panic!("stress agent {i} shutdown: {e}"));
i
}));
}
let mut results = Vec::new();
for h in handles {
results.push(h.await.expect("tokio task panicked"));
}
results.sort_unstable();
assert_eq!(
results,
(0..N).collect::<Vec<_>>(),
"All {N} stress tasks must complete"
);
assert_eq!(
server.post_count(),
N,
"Expected exactly {N} POST requests from stress test"
);
});
}
#[test]
fn double_shutdown_is_safe() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockServer::start("double-shutdown").await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "double"))
.await
.expect("agent");
agent.chat_text("ping").await.expect("chat");
agent.shutdown().await.expect("first shutdown");
let second = agent.shutdown().await;
eprintln!("Second shutdown result: {second:?}");
});
}
#[test]
fn chat_after_shutdown_returns_error() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockServer::start("post-shutdown-chat").await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "post-shutdown"))
.await
.expect("agent");
agent.chat_text("ping").await.expect("pre-shutdown chat");
agent.shutdown().await.expect("shutdown");
let result = agent.chat_text("after shutdown").await;
result.expect_err("Chat after shutdown MUST return Err");
});
}
#[test]
fn concurrent_create_and_shutdown_storm_no_stall() {
let rt = multi_thread_rt();
rt.block_on(async {
const N: usize = 8;
let server = MockServer::start("storm").await;
let url = server.base_url();
let lifecycle = async {
let agents = futures::future::join_all((0..N).map(|i| {
let url = url.clone();
let name = format!("storm-{i}");
async move { BRIDGE.agent(agent_config(&url, &name)).await }
}))
.await
.into_iter()
.enumerate()
.map(|(i, res)| res.unwrap_or_else(|e| panic!("storm agent {i} creation: {e}")))
.collect::<Vec<_>>();
let chats =
futures::future::join_all(agents.iter().map(|a| a.chat_text("storm ping"))).await;
for (i, res) in chats.into_iter().enumerate() {
let text = res.unwrap_or_else(|e| panic!("storm agent {i} chat: {e}"));
assert!(text.contains("mock:storm"), "storm agent {i} got: {text}");
}
let shutdowns = futures::future::join_all(
agents
.into_iter()
.map(|a| async move { a.shutdown().await }),
)
.await;
for (i, res) in shutdowns.into_iter().enumerate() {
res.unwrap_or_else(|e| panic!("storm agent {i} shutdown: {e}"));
}
};
tokio::time::timeout(std::time::Duration::from_mins(1), lifecycle)
.await
.expect("concurrent create/chat/shutdown storm stalled — possible GIL deadlock");
assert_eq!(
server.post_count(),
N,
"Expected exactly {N} POST requests from storm test"
);
});
}