fix(gateway): harden Gemini endpoint routing

This commit is contained in:
MMEXA
2026-05-18 00:53:34 +00:00
parent 81ff375bfd
commit b004a02e4a
20 changed files with 1348 additions and 83 deletions

View File

@@ -241,6 +241,11 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
upstream_is_stream,
request_requires_body_stream_field(body_json, force_body_stream_field),
);
apply_transport_request_body_semantics(
&mut provider_request_body,
transport,
provider_api_format,
);
if let Some(mapping) =
crate::system_features::reasoning_model_directive_mapping_for_api_format_and_model(
state,
@@ -261,6 +266,11 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
upstream_is_stream,
request_requires_body_stream_field(body_json, force_body_stream_field),
);
apply_transport_request_body_semantics(
&mut provider_request_body,
transport,
provider_api_format,
);
}
if let Some(kiro_auth) = kiro_auth.as_ref() {
@@ -368,6 +378,28 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
})
}
fn apply_transport_request_body_semantics(
provider_request_body: &mut Value,
transport: &GatewayProviderTransportSnapshot,
provider_api_format: &str,
) {
if !crate::ai_serving::api_format_alias_matches(provider_api_format, "gemini:embedding")
|| !crate::ai_serving::transport::vertex::is_vertex_transport_context(transport)
{
return;
}
let Some(object) = provider_request_body.as_object_mut() else {
return;
};
if object.contains_key("requests") {
return;
}
object.remove("model");
}
async fn resolve_local_gemini_image_to_openai_image_candidate_payload_parts(
state: &AppState,
parts: &http::request::Parts,

View File

@@ -12,7 +12,7 @@ use crate::control::GatewayControlDecision;
use crate::usage::spawn_sync_report;
use crate::{usage::GatewaySyncReportRequest, AppState, GatewayError};
use axum::body::Body;
use axum::http::Response;
use axum::http::{Response, StatusCode};
use base64::Engine as _;
use tracing::warn;
@@ -148,6 +148,74 @@ fn build_local_core_sync_finalize_fallback_response(
build_local_sync_response_from_bytes(trace_id, decision, payload, Vec::new())
}
fn maybe_build_invalid_provider_success_finalize_response(
trace_id: &str,
decision: &GatewayControlDecision,
payload: &GatewaySyncReportRequest,
) -> Result<Option<Response<Body>>, GatewayError> {
if !local_core_sync_finalize_has_invalid_provider_success(payload)? {
return Ok(None);
}
let client_api_format = resolve_local_sync_client_api_format(payload);
let message = "Provider returned HTTP 200 but the Gemini response did not contain visible model output; refusing to finalize it as a successful response.";
let body_json = build_core_error_body_for_client_format(
&client_api_format,
message,
Some("invalid_provider_success_response"),
LocalCoreSyncErrorKind::ServerError,
)
.unwrap_or_else(|| {
serde_json::json!({
"error": {
"message": message,
"type": "server_error",
"code": "invalid_provider_success_response"
}
})
});
let mut response_headers = payload.headers.clone();
response_headers.remove("content-encoding");
response_headers.remove("content-length");
response_headers.insert("content-type".to_string(), "application/json".to_string());
let body_bytes =
serde_json::to_vec(&body_json).map_err(|err| GatewayError::Internal(err.to_string()))?;
response_headers.insert("content-length".to_string(), body_bytes.len().to_string());
Ok(Some(build_client_response_from_parts(
StatusCode::BAD_GATEWAY.as_u16(),
&response_headers,
Body::from(body_bytes),
trace_id,
Some(decision),
)?))
}
fn local_core_sync_finalize_has_invalid_provider_success(
payload: &GatewaySyncReportRequest,
) -> Result<bool, GatewayError> {
if payload.status_code >= 400 || !is_core_error_finalize_kind(payload.report_kind.as_str()) {
return Ok(false);
}
let provider_api_format = resolve_local_sync_provider_api_format(payload);
if aether_ai_formats::normalize_api_format_alias(&provider_api_format)
!= "gemini:generate_content"
{
return Ok(false);
}
let Some(body_json) = resolve_local_sync_source_body_json(payload)? else {
return Ok(false);
};
if has_nested_error(&body_json) {
return Ok(false);
}
Ok(
aether_ai_formats::formats::gemini::generate_content::response::from_raw(&body_json)
.is_none(),
)
}
pub(crate) fn build_best_effort_local_core_error_body(
payload: &GatewaySyncReportRequest,
body_json: &serde_json::Value,
@@ -283,6 +351,16 @@ fn resolve_local_sync_client_api_format(payload: &GatewaySyncReportRequest) -> S
.to_ascii_lowercase()
}
fn resolve_local_sync_provider_api_format(payload: &GatewaySyncReportRequest) -> String {
payload
.report_context
.as_ref()
.and_then(|value| value.get("provider_api_format"))
.and_then(|value| value.as_str())
.map(|value| value.trim().to_ascii_lowercase())
.unwrap_or_else(|| resolve_local_sync_client_api_format(payload))
}
pub(crate) fn resolve_core_error_background_report_kind(report_kind: &str) -> Option<String> {
core_error_background_report_kind(report_kind).map(ToOwned::to_owned)
}
@@ -522,6 +600,10 @@ pub(crate) async fn submit_local_core_error_or_sync_finalize(
maybe_compile_sync_finalize_response(trace_id, decision, &payload)?
{
response
} else if let Some(response) =
maybe_build_invalid_provider_success_finalize_response(trace_id, decision, &payload)?
{
response
} else if let Some(response) =
maybe_build_local_core_error_response(trace_id, decision, &payload)?
{
@@ -566,9 +648,10 @@ mod tests {
use axum::body::to_bytes;
use serde_json::json;
use super::maybe_build_local_core_error_response;
use super::{maybe_build_local_core_error_response, submit_local_core_error_or_sync_finalize};
use crate::control::GatewayControlDecision;
use crate::usage::GatewaySyncReportRequest;
use crate::AppState;
fn test_decision() -> GatewayControlDecision {
GatewayControlDecision::synthetic(
@@ -684,4 +767,59 @@ mod tests {
})
);
}
#[tokio::test]
async fn local_core_sync_finalize_rejects_gemini_http_200_without_visible_output() {
let mut payload = core_finalize_payload(
"openai_chat_sync_finalize",
"openai:chat",
"gemini:generate_content",
200,
json!({
"candidates": [{
"content": {"role": "model"},
"finishReason": "MAX_TOKENS"
}],
"usageMetadata": {
"promptTokenCount": 8,
"candidatesTokenCount": 1,
"thoughtsTokenCount": 25,
"totalTokenCount": 34
},
"modelVersion": "gemini-3-flash-preview",
"responseId": "resp-empty"
}),
);
payload.report_context = Some(json!({
"client_api_format": "openai:chat",
"provider_api_format": "gemini:generate_content",
"needs_conversion": true,
"has_envelope": false
}));
let state = AppState::new().expect("state should build");
let response = submit_local_core_error_or_sync_finalize(
&state,
"trace-invalid-gemini-200",
&test_decision(),
payload,
)
.await
.expect("response should build");
assert_eq!(response.status(), http::StatusCode::BAD_GATEWAY);
let body: serde_json::Value = serde_json::from_slice(
&to_bytes(response.into_body(), usize::MAX)
.await
.expect("body should read"),
)
.expect("body should decode");
let message = body["error"]["message"]
.as_str()
.expect("error message should exist");
assert!(
message.contains("visible model output"),
"unexpected message: {message}"
);
}
}

View File

@@ -3,7 +3,10 @@ use std::io::Error as IoError;
use std::sync::Arc;
use std::time::{Duration, Instant};
use aether_contracts::{ExecutionPlan, ExecutionResult, ExecutionTelemetry};
use aether_contracts::{
ExecutionError, ExecutionErrorKind, ExecutionPhase, ExecutionPlan, ExecutionResult,
ExecutionTelemetry,
};
use aether_data_contracts::repository::candidates::RequestCandidateStatus;
use aether_scheduler_core::{
execution_error_details, parse_request_candidate_report_context,
@@ -26,8 +29,8 @@ use tokio::time::MissedTickBehavior;
use tracing::{debug, warn};
use crate::ai_serving::api::{
implicit_sync_finalize_report_kind, maybe_build_sync_finalize_outcome,
LocalCoreSyncFinalizeOutcome,
build_core_error_body_for_client_format, implicit_sync_finalize_report_kind,
maybe_build_sync_finalize_outcome, LocalCoreSyncErrorKind, LocalCoreSyncFinalizeOutcome,
};
use crate::api::response::{
attach_control_metadata_headers, build_client_response, build_client_response_from_parts,
@@ -183,6 +186,55 @@ fn build_sync_report_payload(
}
}
fn invalid_gemini_provider_success_message(
plan: &ExecutionPlan,
report_context: Option<&Value>,
status_code: u16,
body_json: Option<&Value>,
) -> Option<&'static str> {
if status_code >= 400 {
return None;
}
let provider_api_format = report_context
.and_then(|value| value.get("provider_api_format"))
.and_then(Value::as_str)
.unwrap_or(plan.provider_api_format.as_str());
if aether_ai_formats::normalize_api_format_alias(provider_api_format)
!= "gemini:generate_content"
{
return None;
}
let body_json = body_json?;
if body_json
.as_object()
.is_some_and(|object| object.get("error").is_some_and(|error| !error.is_null()))
{
return None;
}
if aether_ai_formats::formats::gemini::generate_content::response::from_raw(body_json).is_some()
{
return None;
}
Some("Provider returned HTTP 200 but the Gemini response did not contain visible model output; refusing to finalize it as a successful response.")
}
fn build_invalid_provider_success_body(
plan: &ExecutionPlan,
report_context: Option<&Value>,
message: &str,
) -> Option<Value> {
let client_api_format = report_context
.and_then(|value| value.get("client_api_format"))
.and_then(Value::as_str)
.unwrap_or(plan.client_api_format.as_str());
build_core_error_body_for_client_format(
client_api_format,
message,
Some("invalid_provider_success_response"),
LocalCoreSyncErrorKind::ServerError,
)
}
#[derive(Debug, Clone)]
struct OpenAiImageSyncProgressSnapshot {
phase: &'static str,
@@ -1337,19 +1389,37 @@ async fn execute_execution_runtime_sync_impl(
local_failover_response_text,
local_failover_analysis,
) = loop {
let result_body_json = result
.body
.as_ref()
.and_then(|body| body.json_body.as_ref());
let (result_error_type, result_error_message) =
execution_error_details(result.error.as_ref(), result_body_json);
let result_latency_ms = result
.telemetry
.as_ref()
.and_then(|telemetry| telemetry.elapsed_ms);
let mut headers = std::mem::take(&mut result.headers);
let (body_bytes, body_json, body_base64) =
let (body_bytes, mut body_json, body_base64) =
decode_execution_result_body(result.body.take(), &mut headers)?;
if let Some(message) = invalid_gemini_provider_success_message(
&plan,
report_context.as_ref(),
result.status_code,
body_json.as_ref(),
) {
result.status_code = StatusCode::BAD_GATEWAY.as_u16();
result.error = Some(ExecutionError {
kind: ExecutionErrorKind::Upstream5xx,
phase: ExecutionPhase::Finalize,
message: message.to_string(),
upstream_status: Some(StatusCode::OK.as_u16()),
retryable: false,
failover_recommended: false,
});
if let Some(error_body) =
build_invalid_provider_success_body(&plan, report_context.as_ref(), message)
{
body_json = Some(error_body);
headers.insert("content-type".to_string(), "application/json".to_string());
}
}
let (result_error_type, result_error_message) =
execution_error_details(result.error.as_ref(), body_json.as_ref());
let local_failover_response_text = local_failover_response_text(
body_json.as_ref(),
&body_bytes,
@@ -2058,6 +2128,41 @@ mod tests {
}
}
fn test_gemini_chat_plan() -> ExecutionPlan {
let mut plan = test_openai_image_plan(false);
plan.client_api_format = "openai:chat".to_string();
plan.provider_api_format = "gemini:generate_content".to_string();
plan.model_name = Some("gemini-3-flash-preview".to_string());
plan
}
#[test]
fn invalid_gemini_provider_success_uses_plan_format_when_context_is_missing() {
let plan = test_gemini_chat_plan();
let body = json!({
"candidates": [{
"content": {"role": "model"},
"finishReason": "MAX_TOKENS"
}],
"usageMetadata": {
"promptTokenCount": 8,
"candidatesTokenCount": 1,
"thoughtsTokenCount": 25,
"totalTokenCount": 34
}
});
let message = invalid_gemini_provider_success_message(
&plan,
None,
StatusCode::OK.as_u16(),
Some(&body),
)
.expect("empty Gemini 200 response should be rejected from plan format");
assert!(message.contains("visible model output"));
}
#[tokio::test]
async fn json_whitespace_heartbeat_stream_prefixes_final_json() {
let (tx, rx) = mpsc::channel::<Result<Bytes, IoError>>(1);

View File

@@ -18,6 +18,7 @@ use crate::ai_serving::{
};
use crate::clock::current_unix_ms;
use crate::execution_runtime;
use crate::handlers::admin::provider::write::provider::reconcile_admin_fixed_provider_template_endpoints;
use crate::handlers::admin::request::{AdminAppState, AdminGatewayProviderTransportSnapshot};
use crate::handlers::shared::provider_pool::{
admin_provider_pool_config_from_config_value, read_admin_provider_pool_runtime_state,
@@ -555,7 +556,6 @@ fn provider_query_build_test_request_body_with_model_policy(
"content": provider_query_extract_message(payload)
.unwrap_or_else(|| DEFAULT_PROVIDER_QUERY_TEST_MESSAGE.to_string())
}],
"max_tokens": 30,
"temperature": 0.7,
"stream": true,
})
@@ -1022,6 +1022,8 @@ async fn provider_query_build_kiro_test_candidates(
payload: &Value,
requested_model_override: Option<&str>,
) -> Result<Vec<ProviderQueryTestCandidate>, Response<Body>> {
provider_query_reconcile_fixed_provider_endpoints_for_test_model(state, provider).await?;
let provider_ids = vec![provider.id.clone()];
let endpoints = state
.app()
@@ -1227,6 +1229,33 @@ async fn provider_query_build_kiro_test_candidates(
Ok(candidates)
}
async fn provider_query_reconcile_fixed_provider_endpoints_for_test_model(
state: &AdminAppState<'_>,
provider: &StoredProviderCatalogProvider,
) -> Result<(), Response<Body>> {
if state
.fixed_provider_template(&provider.provider_type)
.is_none()
|| !state.has_provider_catalog_data_writer()
{
return Ok(());
}
reconcile_admin_fixed_provider_template_endpoints(state, provider)
.await
.map_err(|err| {
warn!(
provider_id = %provider.id,
provider_type = %provider.provider_type,
error = ?err,
"admin provider-query test-model: failed to reconcile fixed provider endpoints"
);
build_admin_provider_query_bad_request_response(
ADMIN_PROVIDER_QUERY_NO_ACTIVE_API_KEY_DETAIL,
)
})
}
fn provider_query_decode_execution_body(
result: &aether_contracts::ExecutionResult,
) -> Option<Vec<u8>> {
@@ -1256,7 +1285,7 @@ fn provider_query_standard_execution_response_body(
provider_api_format: &str,
result: &aether_contracts::ExecutionResult,
) -> Option<Value> {
result
let body = result
.body
.as_ref()
.and_then(|body| body.json_body.clone())
@@ -1264,7 +1293,15 @@ fn provider_query_standard_execution_response_body(
provider_query_decode_execution_body(result).and_then(|body| {
provider_query_aggregate_standard_stream_sync_response(provider_api_format, &body)
})
})
})?;
if result.status_code < 400
&& provider_query_normalize_api_format_alias(provider_api_format)
== "gemini:generate_content"
&& aether_ai_formats::formats::gemini::generate_content::response::from_raw(&body).is_none()
{
return None;
}
Some(body)
}
fn provider_query_extract_error_message(

View File

@@ -28,6 +28,17 @@ fn provider_query_test_request_body_defaults_missing_model() {
assert_eq!(body["model"], json!("fallback-model"));
}
#[test]
fn provider_query_default_test_request_body_does_not_set_max_tokens() {
let body = provider_query_build_test_request_body(&json!({}), "fallback-model");
assert_eq!(body["model"], json!("fallback-model"));
assert!(
body.get("max_tokens").is_none(),
"admin model test must not silently force a low max_tokens value"
);
}
#[test]
fn provider_query_failover_request_body_overrides_custom_model() {
let payload = json!({
@@ -168,6 +179,40 @@ fn provider_query_standard_test_aggregates_responses_stream_body() {
assert_eq!(body["output"][0]["content"][0]["text"], json!("Hello"));
}
#[test]
fn provider_query_standard_test_rejects_gemini_success_without_visible_output() {
let result = aether_contracts::ExecutionResult {
request_id: "provider-test".to_string(),
candidate_id: Some("candidate-0".to_string()),
status_code: 200,
headers: BTreeMap::new(),
body: Some(aether_contracts::ResponseBody {
json_body: Some(json!({
"candidates": [{
"content": {"role": "model"},
"finishReason": "MAX_TOKENS"
}],
"usageMetadata": {
"promptTokenCount": 8,
"candidatesTokenCount": 1,
"thoughtsTokenCount": 25,
"totalTokenCount": 34
},
"modelVersion": "gemini-3-flash-preview",
"responseId": "resp-empty"
})),
body_bytes_b64: None,
}),
telemetry: None,
error: None,
};
assert!(
provider_query_standard_execution_response_body("gemini:generate_content", &result)
.is_none()
);
}
#[test]
fn provider_query_test_adapter_routes_fixed_provider_endpoint_types() {
assert_eq!(

View File

@@ -179,9 +179,6 @@ pub(super) async fn maybe_build_local_test_connection_route_response(
"role": "user",
"parts": [{"text": "Health check"}],
}],
"generationConfig": {
"maxOutputTokens": 5,
},
}),
_ => return None,
};

View File

@@ -129,6 +129,56 @@ fn gemini_embedding_success_state(
.with_data_state_for_tests(data_state)
}
fn vertex_gemini_embedding_success_state(execution_runtime_url: String) -> AppState {
let mut snapshot = sample_currently_usable_auth_snapshot(
"key-vertex-gemini-embedding-success",
"user-vertex-gemini-embedding-success",
);
snapshot.user_allowed_providers = None;
snapshot.api_key_allowed_providers = Some(vec!["openai".to_string(), "vertex_ai".to_string()]);
snapshot.user_allowed_api_formats = Some(vec!["openai:embedding".to_string()]);
snapshot.api_key_allowed_api_formats = Some(vec!["openai:embedding".to_string()]);
snapshot.user_allowed_models = Some(vec!["gemini-embedding-2-preview".to_string()]);
snapshot.api_key_allowed_models = Some(vec!["gemini-embedding-2-preview".to_string()]);
let auth_repository = Arc::new(InMemoryAuthApiKeySnapshotRepository::seed(vec![(
Some(hash_api_key("sk-vertex-gemini-embedding-success")),
snapshot,
)]));
let candidate_repository =
Arc::new(InMemoryMinimalCandidateSelectionReadRepository::seed(vec![
vertex_gemini_embedding_candidate_row(),
]));
let mut provider = sample_provider("provider-vertex-gemini-embedding", "Vertex AI", 1);
provider.provider_type = "vertex_ai".to_string();
let mut key = sample_key(
"key-upstream-vertex-gemini-embedding",
"provider-vertex-gemini-embedding",
"gemini:embedding",
"sk-upstream-vertex-gemini-embedding",
);
key.allowed_models = Some(json!(["gemini-embedding-2"]));
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
vec![sample_endpoint(
"endpoint-vertex-gemini-embedding",
"provider-vertex-gemini-embedding",
"gemini:embedding",
"https://aiplatform.googleapis.com",
)],
vec![key],
));
let data_state =
GatewayDataState::with_provider_catalog_and_minimal_candidate_selection_for_tests(
provider_catalog_repository,
candidate_repository,
)
.with_auth_api_key_reader(auth_repository)
.with_encryption_key_for_tests(DEVELOPMENT_ENCRYPTION_KEY);
build_state_with_execution_runtime_override(execution_runtime_url)
.with_data_state_for_tests(data_state)
}
fn gemini_embedding_conversion_execution_runtime() -> Router {
Router::new().route(
"/v1/execute/sync",
@@ -139,6 +189,16 @@ fn gemini_embedding_conversion_execution_runtime() -> Router {
)
}
fn vertex_gemini_embedding_conversion_execution_runtime() -> Router {
Router::new().route(
"/v1/execute/sync",
any(|Json(plan): Json<ExecutionPlan>| async move {
assert_openai_to_vertex_gemini_embedding_execution_plan(&plan);
Json(gemini_embedding_execution_result(&plan))
}),
)
}
fn gemini_embedding_batch_conversion_execution_runtime() -> Router {
Router::new().route(
"/v1/execute/sync",
@@ -227,6 +287,19 @@ fn gemini_embedding_candidate_row() -> StoredMinimalCandidateSelectionRow {
}
}
fn vertex_gemini_embedding_candidate_row() -> StoredMinimalCandidateSelectionRow {
let mut row = gemini_embedding_candidate_row();
row.provider_id = "provider-vertex-gemini-embedding".to_string();
row.provider_name = "Vertex AI".to_string();
row.provider_type = "vertex_ai".to_string();
row.endpoint_id = "endpoint-vertex-gemini-embedding".to_string();
row.key_id = "key-upstream-vertex-gemini-embedding".to_string();
row.key_name = "default".to_string();
row.key_allowed_models = Some(vec!["gemini-embedding-2".to_string()]);
row.model_provider_model_name = "gemini-embedding-2".to_string();
row
}
fn assert_embedding_execution_plan(plan: &ExecutionPlan) {
assert_eq!(plan.client_api_format, "openai:embedding");
assert_eq!(plan.provider_api_format, "openai:embedding");
@@ -262,6 +335,30 @@ fn assert_openai_to_gemini_embedding_execution_plan(plan: &ExecutionPlan) {
assert!(body.get("messages").is_none());
}
fn assert_openai_to_vertex_gemini_embedding_execution_plan(plan: &ExecutionPlan) {
assert_eq!(plan.provider_id, "provider-vertex-gemini-embedding");
assert_eq!(plan.client_api_format, "openai:embedding");
assert_eq!(plan.provider_api_format, "gemini:embedding");
assert_eq!(plan.method, "POST");
assert_eq!(
plan.url,
"https://aiplatform.googleapis.com/v1/publishers/google/models/gemini-embedding-2:embedContent?key=sk-upstream-vertex-gemini-embedding"
);
assert_eq!(
plan.model_name.as_deref(),
Some("gemini-embedding-2-preview")
);
assert!(!plan.stream);
let body = plan.body.json_body.as_ref().expect("json request body");
assert!(
body.get("model").is_none(),
"Vertex embedContent carries the model in the path; the body must not repeat it"
);
assert_eq!(body["content"]["parts"][0]["text"], "hello");
assert!(body.get("input").is_none());
assert!(body.get("messages").is_none());
}
fn assert_openai_to_gemini_batch_embedding_execution_plan(plan: &ExecutionPlan) {
assert_eq!(plan.client_api_format, "openai:embedding");
assert_eq!(plan.provider_api_format, "gemini:embedding");
@@ -502,6 +599,50 @@ async fn embeddings_route_converts_openai_payload_to_gemini_embedding_provider()
execution_runtime_handle.abort();
}
#[tokio::test]
async fn embeddings_route_converts_openai_payload_to_vertex_gemini_embedding_provider() {
let (execution_runtime_url, execution_runtime_handle) =
start_server(vertex_gemini_embedding_conversion_execution_runtime()).await;
let gateway =
build_router_with_state(vertex_gemini_embedding_success_state(execution_runtime_url));
let (gateway_url, gateway_handle) = start_server(gateway).await;
let response = reqwest::Client::new()
.post(format!("{gateway_url}/v1/embeddings"))
.header(
http::header::AUTHORIZATION,
"Bearer sk-vertex-gemini-embedding-success",
)
.json(&json!({
"model": "gemini-embedding-2-preview",
"input": "hello"
}))
.send()
.await
.expect("request should succeed");
let endpoint_signature = response
.headers()
.get(CONTROL_ENDPOINT_SIGNATURE_HEADER)
.and_then(|value| value.to_str().ok())
.map(str::to_string);
let status = response.status();
let body_text = response.text().await.expect("body should read");
assert_eq!(
status,
StatusCode::OK,
"unexpected response body: {body_text}"
);
assert_eq!(endpoint_signature.as_deref(), Some("openai:embedding"));
let payload: serde_json::Value = serde_json::from_str(&body_text).expect("body should parse");
assert_eq!(payload["object"], "list");
assert_eq!(payload["model"], "gemini-embedding-2-preview");
assert_eq!(payload["data"][0]["embedding"], json!([0.1, 0.2, 0.3]));
gateway_handle.abort();
execution_runtime_handle.abort();
}
#[tokio::test]
async fn embeddings_route_converts_openai_batch_payload_to_gemini_batch_endpoint() {
let (execution_runtime_url, execution_runtime_handle) =

View File

@@ -1978,6 +1978,91 @@ async fn gateway_handles_public_test_connection_without_hitting_fallback_probe()
provider_handle.abort();
}
#[tokio::test]
async fn gateway_gemini_test_connection_does_not_force_low_max_output_tokens() {
let provider_hits = Arc::new(Mutex::new(0usize));
let provider_hits_clone = Arc::clone(&provider_hits);
let provider = Router::new().route(
"/{*path}",
any(move |request: Request| {
let provider_hits_inner = Arc::clone(&provider_hits_clone);
async move {
*provider_hits_inner.lock().expect("mutex should lock") += 1;
let body = to_bytes(request.into_body(), usize::MAX)
.await
.expect("body should read");
let body_json: serde_json::Value =
serde_json::from_slice(&body).expect("json body should parse");
assert_eq!(body_json["contents"][0]["parts"][0]["text"], "Health check");
assert!(
body_json
.get("generationConfig")
.and_then(|config| config.get("maxOutputTokens"))
.is_none(),
"Gemini test connection must not force a tiny maxOutputTokens value"
);
Json(json!({
"candidates": [{
"content": {
"role": "model",
"parts": [{"text": "ok"}]
},
"finishReason": "STOP"
}],
"responseId": "gemini_test_connection_ok"
}))
.into_response()
}
}),
);
let (provider_url, provider_handle) = start_server(provider).await;
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![sample_provider("provider-gemini", "google", 10)],
vec![sample_endpoint(
"endpoint-gemini",
"provider-gemini",
"gemini:generate_content",
&provider_url,
)],
vec![sample_key(
"key-gemini",
"provider-gemini",
"gemini:generate_content",
"google-api-key",
)],
));
let gateway = build_router_with_state(
AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(GatewayDataState::with_provider_transport_reader_for_tests(
provider_catalog_repository,
DEVELOPMENT_ENCRYPTION_KEY,
)),
);
let (gateway_url, gateway_handle) = start_server(gateway).await;
let response = reqwest::Client::new()
.get(format!(
"{gateway_url}/v1/test-connection?provider=provider-gemini&model=gemini-3-flash-preview&api_format=gemini:generate_content"
))
.send()
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
assert_eq!(payload["status"], "success");
assert_eq!(payload["provider_id"], "provider-gemini");
assert_eq!(payload["endpoint_id"], "endpoint-gemini");
assert_eq!(payload["api_format"], "gemini:generate_content");
assert_eq!(*provider_hits.lock().expect("mutex should lock"), 1);
gateway_handle.abort();
provider_handle.abort();
}
async fn assert_public_support_route_returns_local_503(
method: reqwest::Method,
path: &str,