use agy_bridge_test_support::*;
#[test]
fn empty_candidate_then_text_recovers() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockGeminiServer::start(vec![
MockResponse::EmptyCandidate,
MockResponse::Text("recovered after thinking-only".into()),
])
.await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "recoverable"))
.await
.expect("agent");
let result = tokio::time::timeout(
std::time::Duration::from_secs(30),
agent.chat_text("trigger empty candidate"),
)
.await;
match result {
Ok(Ok(text)) => {
eprintln!("Agent recovered with text: {text:?}");
assert_eq!(
text.trim(),
"recovered after thinking-only",
"Agent should receive the retry text after empty candidate"
);
}
Ok(Err(e)) => {
let msg = e.to_string();
eprintln!("Agent got error (acceptable if backend doesn't retry): {msg}");
assert!(
msg.contains("model output") || msg.contains("empty"),
"Error should be about model output quality, not a stream abort. Got: {msg}"
);
}
Err(elapsed) => {
panic!("Timed out waiting for agent — possible stream deadlock: {elapsed}");
}
}
assert!(
server.post_count() >= 1,
"Server should have received at least 1 POST"
);
agent.shutdown().await.expect("shutdown");
});
}
#[test]
fn multiple_empty_candidates_then_text_recovers() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockGeminiServer::start(vec![
MockResponse::EmptyCandidate,
MockResponse::EmptyCandidate,
MockResponse::Text("survived two empties".into()),
])
.await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "multi-empty"))
.await
.expect("agent");
let result = tokio::time::timeout(
std::time::Duration::from_mins(1),
agent.chat_text("trigger multiple empties"),
)
.await;
match result {
Ok(Ok(text)) => {
eprintln!("Recovered after empty candidates: {text:?}");
let trimmed = text.trim();
assert!(
trimmed == "survived two empties" || trimmed.is_empty(),
"Expected retry text or empty, got: {trimmed:?}"
);
}
Ok(Err(e)) => {
let msg = e.to_string();
eprintln!("Error after multiple empties (acceptable): {msg}");
assert!(
!msg.contains("channel closed") && !msg.contains("ChannelClosed"),
"Stream must not abort on recoverable errors. Got: {msg}"
);
}
Err(elapsed) => {
panic!("Timed out — stream may be stuck after repeated empties: {elapsed}");
}
}
assert!(
server.post_count() >= 1,
"Server should have been called at least once"
);
agent.shutdown().await.expect("shutdown");
});
}
#[test]
fn streaming_handle_survives_empty_candidate() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockGeminiServer::start(vec![
MockResponse::EmptyCandidate,
MockResponse::Text("stream recovered".into()),
])
.await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "stream-recover"))
.await
.expect("agent");
let handle_result = tokio::time::timeout(
std::time::Duration::from_secs(30),
agent.chat("trigger empty via handle"),
)
.await;
match handle_result {
Ok(Ok(handle)) => {
let text_result =
tokio::time::timeout(std::time::Duration::from_secs(30), handle.text()).await;
match text_result {
Ok(Ok(chat_result)) => {
let text = chat_result.text();
eprintln!("Streaming handle recovered: {text:?}");
assert_eq!(
text.trim(),
"stream recovered",
"Streaming handle should deliver retry text"
);
}
Ok(Err(e)) => {
let msg = e.to_string();
eprintln!("Streaming handle error (acceptable): {msg}");
assert!(
msg.contains("model output") || msg.contains("empty"),
"Error should be model-quality, not stream abort. Got: {msg}"
);
}
Err(elapsed) => {
panic!("Streaming handle timed out — possible deadlock: {elapsed}");
}
}
}
Ok(Err(e)) => {
eprintln!("chat() error (acceptable if backend doesn't retry): {e}");
}
Err(elapsed) => {
panic!("Timed out at chat() level: {elapsed}");
}
}
agent.shutdown().await.expect("shutdown");
});
}
#[test]
fn healthy_agent_unaffected_by_sibling_empty_candidate() {
let rt = multi_thread_rt();
rt.block_on(async {
let healthy_server =
MockGeminiServer::start(vec![MockResponse::Text("healthy response".into())]).await;
let empty_then_text_server = MockGeminiServer::start(vec![
MockResponse::EmptyCandidate,
MockResponse::Text("recovered".into()),
])
.await;
let healthy_agent = BRIDGE
.agent(agent_config(&healthy_server.base_url(), "healthy"))
.await
.expect("healthy agent");
let recovering_agent = BRIDGE
.agent(agent_config(
&empty_then_text_server.base_url(),
"recovering",
))
.await
.expect("recovering agent");
let (healthy_res, recovering_res) = tokio::join!(
healthy_agent.chat_text("ping"),
tokio::time::timeout(
std::time::Duration::from_secs(30),
recovering_agent.chat_text("trigger empty"),
),
);
let healthy_text = healthy_res.expect("healthy agent must succeed");
assert_eq!(healthy_text.trim(), "healthy response");
match recovering_res {
Ok(Ok(text)) => {
eprintln!("Recovering agent got: {text:?}");
assert_eq!(text.trim(), "recovered");
}
Ok(Err(e)) => {
eprintln!("Recovering agent error (acceptable): {e}");
}
Err(elapsed) => {
panic!("Recovering agent timed out — possible stream deadlock: {elapsed}");
}
}
healthy_agent.shutdown().await.expect("shutdown healthy");
recovering_agent
.shutdown()
.await
.expect("shutdown recovering");
});
}
#[test]
fn http_503_returns_err() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockGeminiServer::start(vec![MockResponse::HttpError {
status: 503,
message: "Service unavailable".into(),
}])
.await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "test-503"))
.await
.expect("agent");
let result = tokio::time::timeout(
std::time::Duration::from_secs(30),
agent.chat_text("trigger 503"),
)
.await;
match result {
Ok(Ok(text)) => {
panic!("503 must return Err, got Ok({text:?})");
}
Ok(Err(e)) => {
let msg = e.to_string();
eprintln!("503 error (expected): {msg}");
assert!(
msg.contains("503") || msg.contains("unavailable") || msg.contains("error"),
"Error message should reference the 503: {msg}"
);
}
Err(elapsed) => {
panic!("Timed out waiting for 503 error: {elapsed}");
}
}
agent.shutdown().await.expect("shutdown");
});
}
#[test]
fn http_503_then_recovery_returns_ok() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockGeminiServer::start(vec![
MockResponse::HttpError {
status: 503,
message: "Temporary unavailable".into(),
},
MockResponse::Text("recovered after 503".into()),
])
.await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "test-503-recovery"))
.await
.expect("agent");
let result = tokio::time::timeout(
std::time::Duration::from_secs(30),
agent.chat_text("trigger 503 then recover"),
)
.await;
match result {
Ok(Ok(text)) => {
eprintln!("Recovered from 503: {text:?}");
assert_eq!(
text.trim(),
"recovered after 503",
"Should get the retry text after 503 recovery"
);
}
Ok(Err(e)) => {
let msg = e.to_string();
eprintln!("503 recovery failed (may be SDK version dependent): {msg}");
assert!(
!msg.contains("channel closed") && !msg.contains("ChannelClosed"),
"Must not be a stream abort. Got: {msg}"
);
}
Err(elapsed) => {
panic!("Timed out waiting for 503 recovery: {elapsed}");
}
}
agent.shutdown().await.expect("shutdown");
});
}
#[test]
fn http_429_returns_err() {
let rt = multi_thread_rt();
rt.block_on(async {
let server = MockGeminiServer::start(vec![MockResponse::HttpError {
status: 429,
message: "Quota exceeded".into(),
}])
.await;
let agent = BRIDGE
.agent(agent_config(&server.base_url(), "test-429"))
.await
.expect("agent");
let result = tokio::time::timeout(
std::time::Duration::from_secs(30),
agent.chat_text("trigger 429"),
)
.await;
match result {
Ok(Ok(text)) => {
panic!("429 must return Err, got Ok({text:?})");
}
Ok(Err(e)) => {
let msg = e.to_string();
eprintln!("429 error (expected): {msg}");
assert!(
msg.contains("429") || msg.contains("quota") || msg.contains("error"),
"Error message should reference the 429: {msg}"
);
}
Err(elapsed) => {
panic!("Timed out waiting for 429 error: {elapsed}");
}
}
agent.shutdown().await.expect("shutdown");
});
}