diff --git a/apps/aether-gateway/src/ai_serving/api.rs b/apps/aether-gateway/src/ai_serving/api.rs index 5eadbcf78..282d680a6 100644 --- a/apps/aether-gateway/src/ai_serving/api.rs +++ b/apps/aether-gateway/src/ai_serving/api.rs @@ -49,14 +49,15 @@ pub(crate) use aether_ai_formats::api::{ resolve_claude_sync_spec, resolve_gemini_stream_spec, resolve_gemini_sync_spec, resolve_local_image_stream_spec, resolve_local_image_sync_spec, resolve_local_same_format_stream_spec, resolve_local_same_format_sync_spec, - AiControlPlanRequest, ExecutionRuntimeAuthContext, LocalCoreSyncErrorKind, - LocalOpenAiImageSpec, LocalSameFormatProviderFamily, LocalSameFormatProviderSpec, - LocalStandardSourceFamily, LocalStandardSourceMode, LocalStandardSpec, - StreamingStandardTerminalObserver, EXECUTION_RUNTIME_STREAM_DECISION_ACTION, - EXECUTION_RUNTIME_SYNC_DECISION_ACTION, GEMINI_FILES_DOWNLOAD_PLAN_KIND, - GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, - OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND, - OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, + resolve_openai_embedding_sync_spec, AiControlPlanRequest, ExecutionRuntimeAuthContext, + LocalCoreSyncErrorKind, LocalOpenAiImageSpec, LocalSameFormatProviderFamily, + LocalSameFormatProviderSpec, LocalStandardSourceFamily, LocalStandardSourceMode, + LocalStandardSpec, StreamingStandardTerminalObserver, EXECUTION_RUNTIME_STREAM_DECISION_ACTION, + EXECUTION_RUNTIME_SYNC_DECISION_ACTION, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, + OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, + OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, + OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, OPENAI_VIDEO_CONTENT_PLAN_KIND, OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND, OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND, }; diff --git a/apps/aether-gateway/src/ai_serving/mod.rs b/apps/aether-gateway/src/ai_serving/mod.rs index b64f1ef13..380111930 100644 --- a/apps/aether-gateway/src/ai_serving/mod.rs +++ b/apps/aether-gateway/src/ai_serving/mod.rs @@ -64,7 +64,8 @@ pub(crate) use self::transport::{ candidate_common_transport_skip_reason, candidate_transport_pair_skip_reason, request_conversion_direct_auth, request_conversion_enabled_for_transport, request_conversion_transport_supported, request_conversion_transport_unsupported_reason, - request_pair_allowed_for_transport, CandidateTransportPolicyFacts, + request_pair_allowed_for_transport, request_pair_direct_auth, + request_pair_transport_unsupported_reason, CandidateTransportPolicyFacts, }; pub(crate) use crate::control::GatewayControlDecision; pub(crate) use crate::execution_runtime::{ConversionMode, ExecutionStrategy}; @@ -97,6 +98,28 @@ pub(crate) fn build_provider_transport_request_url( ) } +pub(crate) fn build_provider_transport_request_url_for_request_body( + transport: &GatewayProviderTransportSnapshot, + provider_api_format: &str, + mapped_model: Option<&str>, + upstream_is_stream: bool, + request_query: Option<&str>, + kiro_api_region: Option<&str>, + provider_request_body: Option<&serde_json::Value>, +) -> Option { + self::transport::build_transport_request_url_for_request_body( + transport, + self::transport::TransportRequestUrlParams { + provider_api_format, + mapped_model, + upstream_is_stream, + request_query, + kiro_api_region, + }, + provider_request_body, + ) +} + pub(crate) async fn resolve_execution_runtime_auth_context( state: &AppState, decision: &GatewayControlDecision, diff --git a/apps/aether-gateway/src/ai_serving/planner/common.rs b/apps/aether-gateway/src/ai_serving/planner/common.rs index ae7741196..91cab0f22 100644 --- a/apps/aether-gateway/src/ai_serving/planner/common.rs +++ b/apps/aether-gateway/src/ai_serving/planner/common.rs @@ -13,8 +13,9 @@ pub(crate) use crate::ai_serving::{ EXECUTION_RUNTIME_STREAM_DECISION_ACTION, EXECUTION_RUNTIME_SYNC_ACTION, EXECUTION_RUNTIME_SYNC_DECISION_ACTION, GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, - GEMINI_FILES_DELETE_PLAN_KIND, GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, - GEMINI_FILES_LIST_PLAN_KIND, GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, + GEMINI_EMBEDDING_SYNC_PLAN_KIND, GEMINI_FILES_DELETE_PLAN_KIND, + GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, + GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, diff --git a/apps/aether-gateway/src/ai_serving/planner/decision/control_plan.rs b/apps/aether-gateway/src/ai_serving/planner/decision/control_plan.rs index 7b09b7bba..f681f8558 100644 --- a/apps/aether-gateway/src/ai_serving/planner/decision/control_plan.rs +++ b/apps/aether-gateway/src/ai_serving/planner/decision/control_plan.rs @@ -1,16 +1,16 @@ use crate::ai_serving::planner::common::{ CLAUDE_CHAT_STREAM_PLAN_KIND, CLAUDE_CHAT_SYNC_PLAN_KIND, CLAUDE_CLI_STREAM_PLAN_KIND, CLAUDE_CLI_SYNC_PLAN_KIND, GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, - GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_FILES_DELETE_PLAN_KIND, - GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, - GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, - OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, - OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, - OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND, - OPENAI_RESPONSES_STREAM_PLAN_KIND, OPENAI_RESPONSES_SYNC_PLAN_KIND, - OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, OPENAI_VIDEO_CONTENT_PLAN_KIND, - OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND, - OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND, + GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_FILES_DELETE_PLAN_KIND, GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, + GEMINI_FILES_LIST_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, + GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, + OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, + OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, + OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND, OPENAI_RESPONSES_STREAM_PLAN_KIND, + OPENAI_RESPONSES_SYNC_PLAN_KIND, OPENAI_VIDEO_CANCEL_SYNC_PLAN_KIND, + OPENAI_VIDEO_CONTENT_PLAN_KIND, OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND, + OPENAI_VIDEO_DELETE_SYNC_PLAN_KIND, OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND, }; use crate::ai_serving::planner::plan_builders::{ build_gemini_stream_plan_from_decision, build_gemini_sync_plan_from_decision, @@ -110,7 +110,9 @@ fn build_sync_plan_payload_from_decision( | OPENAI_RERANK_SYNC_PLAN_KIND => { build_standard_sync_plan_from_decision(parts, body_json, payload)? } - GEMINI_CHAT_SYNC_PLAN_KIND | GEMINI_CLI_SYNC_PLAN_KIND => { + GEMINI_CHAT_SYNC_PLAN_KIND + | GEMINI_CLI_SYNC_PLAN_KIND + | GEMINI_EMBEDDING_SYNC_PLAN_KIND => { build_gemini_sync_plan_from_decision(parts, body_json, payload)? } OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND diff --git a/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/family/request.rs b/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/family/request.rs index a2e38b03f..c60accef9 100644 --- a/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/family/request.rs +++ b/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/family/request.rs @@ -254,6 +254,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts( spec, prepared.upstream_is_stream, prepared.kiro_auth.as_ref(), + Some(&provider_request_body), ) else { mark_skipped_local_same_format_provider_candidate_with_failure_diagnostic( state, diff --git a/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/request/url.rs b/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/request/url.rs index 6fc84ae83..b9f064f24 100644 --- a/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/request/url.rs +++ b/apps/aether-gateway/src/ai_serving/planner/passthrough/provider/request/url.rs @@ -14,6 +14,7 @@ pub(crate) fn build_same_format_upstream_url( spec: LocalSameFormatProviderSpec, upstream_is_stream: bool, kiro_auth: Option<&crate::ai_serving::transport::kiro::KiroRequestAuth>, + provider_request_body: Option<&serde_json::Value>, ) -> Option { build_same_format_provider_upstream_url_impl( transport, @@ -23,6 +24,7 @@ pub(crate) fn build_same_format_upstream_url( upstream_is_stream, request_query: parts.uri.query(), kiro_api_region: kiro_auth.map(|auth| auth.auth_config.effective_api_region()), + provider_request_body, }, ) } diff --git a/apps/aether-gateway/src/ai_serving/planner/specialized/image/request.rs b/apps/aether-gateway/src/ai_serving/planner/specialized/image/request.rs index 1af9e5b1c..ea91adf4b 100644 --- a/apps/aether-gateway/src/ai_serving/planner/specialized/image/request.rs +++ b/apps/aether-gateway/src/ai_serving/planner/specialized/image/request.rs @@ -361,6 +361,7 @@ async fn resolve_local_openai_image_to_gemini_candidate_payload_parts( &converted.mapped_model, provider_api_format, upstream_is_stream, + Some(&converted.body_json), ) else { mark_skipped_local_openai_image_candidate_with_failure_diagnostic( state, diff --git a/apps/aether-gateway/src/ai_serving/planner/standard/family/request.rs b/apps/aether-gateway/src/ai_serving/planner/standard/family/request.rs index 1971b0fd7..026924eeb 100644 --- a/apps/aether-gateway/src/ai_serving/planner/standard/family/request.rs +++ b/apps/aether-gateway/src/ai_serving/planner/standard/family/request.rs @@ -75,15 +75,18 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts( .await; } let is_kiro_claude_cli = is_kiro_claude_messages_transport(transport, provider_api_format); - let Some(conversion_kind) = - crate::ai_serving::request_conversion_kind(spec_metadata.api_format, provider_api_format) - else { - return None; - }; - - if let Some(skip_reason) = crate::ai_serving::request_conversion_transport_unsupported_reason( + if !crate::ai_serving::request_pair_allowed_for_transport( transport, - conversion_kind, + spec_metadata.api_format, + provider_api_format, + ) { + return None; + } + + if let Some(skip_reason) = crate::ai_serving::request_pair_transport_unsupported_reason( + transport, + spec_metadata.api_format, + provider_api_format, ) { mark_skipped_local_standard_candidate( state, @@ -156,7 +159,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts( planner_state, transport, candidate, - crate::ai_serving::request_conversion_direct_auth(transport, conversion_kind), + crate::ai_serving::request_pair_direct_auth(transport, provider_api_format), oauth_context, ) .await @@ -286,6 +289,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts( &prepared_candidate.mapped_model, provider_api_format, upstream_is_stream, + Some(&provider_request_body), ) { Some(url) => url, None => { diff --git a/apps/aether-gateway/src/ai_serving/planner/standard/mod.rs b/apps/aether-gateway/src/ai_serving/planner/standard/mod.rs index 7fa8fbad0..2fa17a602 100644 --- a/apps/aether-gateway/src/ai_serving/planner/standard/mod.rs +++ b/apps/aether-gateway/src/ai_serving/planner/standard/mod.rs @@ -38,6 +38,7 @@ pub(crate) use self::openai::{ map_openai_reasoning_effort_to_gemini_budget, maybe_build_stream_local_decision_payload, maybe_build_stream_local_openai_responses_decision_payload, maybe_build_sync_local_decision_payload, + maybe_build_sync_local_openai_embedding_decision_payload, maybe_build_sync_local_openai_responses_decision_payload, parse_openai_stop_sequences, resolve_openai_chat_max_tokens, set_local_openai_chat_execution_exhausted_diagnostic, value_as_u64, @@ -66,14 +67,16 @@ pub(crate) fn build_standard_upstream_url( mapped_model: &str, provider_api_format: &str, upstream_is_stream: bool, + provider_request_body: Option<&serde_json::Value>, ) -> Option { - crate::ai_serving::build_provider_transport_request_url( + crate::ai_serving::build_provider_transport_request_url_for_request_body( transport, provider_api_format, Some(mapped_model), upstream_is_stream, parts.uri.query(), None, + provider_request_body, ) } @@ -85,6 +88,14 @@ pub(crate) async fn maybe_build_sync_local_standard_decision_payload( body_json: &serde_json::Value, plan_kind: &str, ) -> Result, GatewayError> { + if let Some(payload) = self::openai::maybe_build_sync_local_openai_embedding_decision_payload( + state, parts, trace_id, decision, body_json, plan_kind, + ) + .await? + { + return Ok(Some(payload)); + } + if let Some(payload) = self::claude::maybe_build_sync_local_claude_decision_payload( state, parts, trace_id, decision, body_json, plan_kind, ) diff --git a/apps/aether-gateway/src/ai_serving/planner/standard/openai/embedding.rs b/apps/aether-gateway/src/ai_serving/planner/standard/openai/embedding.rs new file mode 100644 index 000000000..490433cb8 --- /dev/null +++ b/apps/aether-gateway/src/ai_serving/planner/standard/openai/embedding.rs @@ -0,0 +1,31 @@ +use crate::ai_serving::resolve_openai_embedding_sync_spec as resolve_surface_sync_spec; +use crate::ai_serving::GatewayControlDecision; +use crate::{AiExecutionDecision, AppState, GatewayError}; + +use super::super::family::maybe_build_sync_via_standard_family_payload; + +pub(crate) fn resolve_sync_spec( + plan_kind: &str, +) -> Option { + resolve_surface_sync_spec(plan_kind) +} + +pub(crate) async fn maybe_build_sync_local_openai_embedding_decision_payload( + state: &AppState, + parts: &http::request::Parts, + trace_id: &str, + decision: &GatewayControlDecision, + body_json: &serde_json::Value, + plan_kind: &str, +) -> Result, GatewayError> { + maybe_build_sync_via_standard_family_payload( + state, + parts, + trace_id, + decision, + body_json, + plan_kind, + resolve_sync_spec, + ) + .await +} diff --git a/apps/aether-gateway/src/ai_serving/planner/standard/openai/mod.rs b/apps/aether-gateway/src/ai_serving/planner/standard/openai/mod.rs index e6456be94..4d10a3aa3 100644 --- a/apps/aether-gateway/src/ai_serving/planner/standard/openai/mod.rs +++ b/apps/aether-gateway/src/ai_serving/planner/standard/openai/mod.rs @@ -1,4 +1,5 @@ mod chat; +mod embedding; mod responses; pub(crate) use crate::ai_serving::{ @@ -14,6 +15,7 @@ pub(crate) use chat::{ maybe_build_stream_local_decision_payload, maybe_build_sync_local_decision_payload, set_local_openai_chat_execution_exhausted_diagnostic, }; +pub(crate) use embedding::maybe_build_sync_local_openai_embedding_decision_payload; pub(crate) use responses::{ build_local_openai_responses_stream_attempt_source_for_kind, build_local_openai_responses_stream_plan_and_reports_for_kind, diff --git a/apps/aether-gateway/src/ai_serving/pure/mod.rs b/apps/aether-gateway/src/ai_serving/pure/mod.rs index bdb09e131..4b367e4d3 100644 --- a/apps/aether-gateway/src/ai_serving/pure/mod.rs +++ b/apps/aether-gateway/src/ai_serving/pure/mod.rs @@ -80,8 +80,8 @@ pub(crate) use aether_ai_formats::api::{ resolve_local_image_stream_spec, resolve_local_image_sync_spec, resolve_local_same_format_stream_spec, resolve_local_same_format_sync_spec, resolve_local_video_sync_spec, resolve_openai_chat_max_tokens, - resolve_openai_responses_stream_spec, resolve_openai_responses_sync_spec, - resolve_requested_gemini_image_model_for_request, + resolve_openai_embedding_sync_spec, resolve_openai_responses_stream_spec, + resolve_openai_responses_sync_spec, resolve_requested_gemini_image_model_for_request, resolve_requested_openai_image_model_for_request, resolve_upstream_is_stream_from_endpoint_config, sanitize_request_path, sanitize_request_path_and_query, sanitize_request_query_string, @@ -119,6 +119,7 @@ pub(crate) use aether_ai_formats::api::{ GEMINI_CLI_STREAM_SUCCESS_REPORT_KIND, GEMINI_CLI_SYNC_ERROR_REPORT_KIND, GEMINI_CLI_SYNC_FINALIZE_REPORT_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND, GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME, + GEMINI_EMBEDDING_SYNC_PLAN_KIND, GEMINI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, GEMINI_FILES_DELETE_PLAN_KIND, GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, GEMINI_VIDEO_CREATE_SYNC_FINALIZE_REPORT_KIND, GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, diff --git a/apps/aether-gateway/src/ai_serving/transport.rs b/apps/aether-gateway/src/ai_serving/transport.rs index 43d074ce9..a9033e64f 100644 --- a/apps/aether-gateway/src/ai_serving/transport.rs +++ b/apps/aether-gateway/src/ai_serving/transport.rs @@ -61,18 +61,19 @@ pub(crate) use aether_provider_transport::{ build_same_format_provider_upstream_url, build_standard_plan_fallback_headers, build_standard_plan_fallback_openai_chat_url, build_standard_plan_fallback_openai_responses_url, build_standard_provider_request_headers, - build_transport_request_url, build_video_create_headers, build_video_create_request_body, - build_video_create_upstream_url, candidate_common_transport_skip_reason, - candidate_transport_pair_skip_reason, classify_same_format_provider_request_behavior, - ensure_upstream_auth_header, gemini_files_transport_unsupported_reason, - header_rules_are_locally_supported, header_rules_have_enabled_rules, - local_gemini_transport_unsupported_reason_with_network, + build_transport_request_url, build_transport_request_url_for_request_body, + build_video_create_headers, build_video_create_request_body, build_video_create_upstream_url, + candidate_common_transport_skip_reason, candidate_transport_pair_skip_reason, + classify_same_format_provider_request_behavior, ensure_upstream_auth_header, + gemini_files_transport_unsupported_reason, header_rules_are_locally_supported, + header_rules_have_enabled_rules, local_gemini_transport_unsupported_reason_with_network, local_openai_chat_transport_unsupported_reason, local_standard_transport_unsupported_reason_with_network, openai_image_transport_unsupported_reason, request_conversion_direct_auth, request_conversion_enabled_for_transport, request_conversion_transport_supported, request_conversion_transport_unsupported_reason, request_pair_allowed_for_transport, - resolve_gemini_files_auth, resolve_openai_image_auth, resolve_same_format_provider_direct_auth, + request_pair_direct_auth, request_pair_transport_unsupported_reason, resolve_gemini_files_auth, + resolve_openai_image_auth, resolve_same_format_provider_direct_auth, resolve_transport_execution_timeouts, resolve_transport_profile, resolve_transport_proxy_snapshot, resolve_transport_proxy_snapshot_with_tunnel_affinity, resolve_video_create_auth, same_format_provider_transport_supported, diff --git a/apps/aether-gateway/src/control/route/ai.rs b/apps/aether-gateway/src/control/route/ai.rs index 4dc2124ba..fa6889032 100644 --- a/apps/aether-gateway/src/control/route/ai.rs +++ b/apps/aether-gateway/src/control/route/ai.rs @@ -104,6 +104,17 @@ pub(super) fn classify_ai_public_route( "gemini:video", true, )) + } else if normalized_path.ends_with(":embedContent") + || normalized_path.ends_with(":batchEmbedContents") + { + Some(classified_with_request_auth_channel( + "ai_public", + "gemini", + "embedding", + "api_key", + "gemini:embedding", + true, + )) } else if is_gemini_cli_request(headers) { Some(classified_with_request_auth_channel( "ai_public", diff --git a/apps/aether-gateway/src/control/route/mod.rs b/apps/aether-gateway/src/control/route/mod.rs index 893e20a28..0e12394db 100644 --- a/apps/aether-gateway/src/control/route/mod.rs +++ b/apps/aether-gateway/src/control/route/mod.rs @@ -229,6 +229,8 @@ pub(super) fn is_gemini_models_route(path: &str) -> bool { (path.starts_with("/v1/models/") || path.starts_with("/v1beta/models/")) && (path.contains(":generateContent") || path.contains(":streamGenerateContent") + || path.contains(":embedContent") + || path.contains(":batchEmbedContents") || path.contains(":predictLongRunning")) } diff --git a/apps/aether-gateway/src/control/tests/ai.rs b/apps/aether-gateway/src/control/tests/ai.rs index e170b593b..731490936 100644 --- a/apps/aether-gateway/src/control/tests/ai.rs +++ b/apps/aether-gateway/src/control/tests/ai.rs @@ -197,6 +197,44 @@ fn classifies_gemini_generate_content_api_key_without_cli_marker() { assert!(decision.is_execution_runtime_candidate()); } +#[test] +fn classifies_gemini_embed_content_as_embedding_route() { + let headers = headers(&[("x-goog-api-key", "gemini-key")]); + let uri: Uri = "/v1beta/models/gemini-embedding-2-preview:embedContent" + .parse() + .expect("uri should parse"); + let decision = + classify_control_route(&http::Method::POST, &uri, &headers).expect("route should classify"); + + assert_eq!(decision.route_family.as_deref(), Some("gemini")); + assert_eq!(decision.route_kind.as_deref(), Some("embedding")); + assert_eq!(decision.request_auth_channel.as_deref(), Some("api_key")); + assert_eq!( + decision.auth_endpoint_signature.as_deref(), + Some("gemini:embedding") + ); + assert!(decision.is_execution_runtime_candidate()); +} + +#[test] +fn classifies_gemini_batch_embed_contents_as_embedding_route() { + let headers = headers(&[("x-goog-api-key", "gemini-key")]); + let uri: Uri = "/v1beta/models/gemini-embedding-2-preview:batchEmbedContents" + .parse() + .expect("uri should parse"); + let decision = + classify_control_route(&http::Method::POST, &uri, &headers).expect("route should classify"); + + assert_eq!(decision.route_family.as_deref(), Some("gemini")); + assert_eq!(decision.route_kind.as_deref(), Some("embedding")); + assert_eq!(decision.request_auth_channel.as_deref(), Some("api_key")); + assert_eq!( + decision.auth_endpoint_signature.as_deref(), + Some("gemini:embedding") + ); + assert!(decision.is_execution_runtime_candidate()); +} + #[test] fn classifies_gemini_predict_long_running_as_video_route() { let headers = headers(&[]); diff --git a/apps/aether-gateway/src/tests/control/proxy/embeddings.rs b/apps/aether-gateway/src/tests/control/proxy/embeddings.rs index 14dedb108..fd68f799c 100644 --- a/apps/aether-gateway/src/tests/control/proxy/embeddings.rs +++ b/apps/aether-gateway/src/tests/control/proxy/embeddings.rs @@ -78,6 +78,87 @@ fn embedding_execution_runtime() -> Router { ) } +fn gemini_embedding_success_state( + execution_runtime_url: String, + client_api_format: &str, +) -> AppState { + let mut snapshot = sample_currently_usable_auth_snapshot( + "key-gemini-embedding-success", + "user-gemini-embedding-success", + ); + snapshot.user_allowed_providers = None; + snapshot.api_key_allowed_providers = None; + snapshot.user_allowed_api_formats = Some(vec![client_api_format.to_string()]); + snapshot.api_key_allowed_api_formats = Some(vec![client_api_format.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-gemini-embedding-success")), + snapshot, + )])); + let candidate_repository = + Arc::new(InMemoryMinimalCandidateSelectionReadRepository::seed(vec![ + gemini_embedding_candidate_row(), + ])); + let mut provider = sample_provider("provider-gemini-embedding", "Gemini Embeddings", 1); + provider.provider_type = "gemini".to_string(); + let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed( + vec![provider], + vec![sample_endpoint( + "endpoint-gemini-embedding", + "provider-gemini-embedding", + "gemini:embedding", + "https://generativelanguage.googleapis.com/v1beta", + )], + vec![sample_key( + "key-upstream-gemini-embedding", + "provider-gemini-embedding", + "gemini:embedding", + "sk-upstream-gemini-embedding", + )], + )); + 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", + any(|Json(plan): Json| async move { + assert_openai_to_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", + any(|Json(plan): Json| async move { + assert_openai_to_gemini_batch_embedding_execution_plan(&plan); + Json(gemini_batch_embedding_execution_result(&plan)) + }), + ) +} + +fn gemini_embedding_native_execution_runtime() -> Router { + Router::new().route( + "/v1/execute/sync", + any(|Json(plan): Json| async move { + assert_native_gemini_embedding_execution_plan(&plan); + Json(gemini_embedding_execution_result(&plan)) + }), + ) +} + fn embedding_candidate_row() -> StoredMinimalCandidateSelectionRow { StoredMinimalCandidateSelectionRow { provider_id: "provider-embedding".to_string(), @@ -112,6 +193,40 @@ fn embedding_candidate_row() -> StoredMinimalCandidateSelectionRow { } } +fn gemini_embedding_candidate_row() -> StoredMinimalCandidateSelectionRow { + StoredMinimalCandidateSelectionRow { + provider_id: "provider-gemini-embedding".to_string(), + provider_name: "Gemini Embeddings".to_string(), + provider_type: "gemini".to_string(), + provider_priority: 1, + provider_is_active: true, + endpoint_id: "endpoint-gemini-embedding".to_string(), + endpoint_api_format: "gemini:embedding".to_string(), + endpoint_api_family: Some("gemini".to_string()), + endpoint_kind: Some("embedding".to_string()), + endpoint_is_active: true, + key_id: "key-upstream-gemini-embedding".to_string(), + key_name: "default".to_string(), + key_auth_type: "api_key".to_string(), + key_is_active: true, + key_api_formats: Some(vec!["gemini:embedding".to_string()]), + key_allowed_models: None, + key_capabilities: None, + key_internal_priority: 50, + key_global_priority_by_format: None, + model_id: "model-gemini-embedding-preview".to_string(), + global_model_id: "global-gemini-embedding-preview".to_string(), + global_model_name: "gemini-embedding-2-preview".to_string(), + global_model_mappings: None, + global_model_supports_streaming: Some(false), + model_provider_model_name: "gemini-embedding-2-preview".to_string(), + model_provider_model_mappings: None, + model_supports_streaming: Some(false), + model_is_active: true, + model_is_available: true, + } +} + fn assert_embedding_execution_plan(plan: &ExecutionPlan) { assert_eq!(plan.client_api_format, "openai:embedding"); assert_eq!(plan.provider_api_format, "openai:embedding"); @@ -123,6 +238,78 @@ fn assert_embedding_execution_plan(plan: &ExecutionPlan) { assert!(body.get("input").is_some()); } +fn assert_openai_to_gemini_embedding_execution_plan(plan: &ExecutionPlan) { + 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://generativelanguage.googleapis.com/v1beta/models/gemini-embedding-2-preview:embedContent" + ); + assert_eq!( + plan.headers.get("x-goog-api-key").map(String::as_str), + Some("sk-upstream-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_eq!(body["model"], "gemini-embedding-2-preview"); + 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"); + assert_eq!(plan.method, "POST"); + assert_eq!( + plan.url, + "https://generativelanguage.googleapis.com/v1beta/models/gemini-embedding-2-preview:batchEmbedContents" + ); + assert_eq!( + plan.headers.get("x-goog-api-key").map(String::as_str), + Some("sk-upstream-gemini-embedding") + ); + assert!(!plan.stream); + let body = plan.body.json_body.as_ref().expect("json request body"); + assert!(body.get("model").is_none()); + let requests = body["requests"].as_array().expect("batch requests"); + assert_eq!(requests.len(), 2); + assert_eq!(requests[0]["model"], "models/gemini-embedding-2-preview"); + assert_eq!(requests[0]["content"]["parts"][0]["text"], "hello"); + assert_eq!(requests[1]["model"], "models/gemini-embedding-2-preview"); + assert_eq!(requests[1]["content"]["parts"][0]["text"], "world"); + assert!(body.get("input").is_none()); + assert!(body.get("messages").is_none()); +} + +fn assert_native_gemini_embedding_execution_plan(plan: &ExecutionPlan) { + assert_eq!(plan.client_api_format, "gemini:embedding"); + assert_eq!(plan.provider_api_format, "gemini:embedding"); + assert_eq!(plan.method, "POST"); + assert_eq!( + plan.url, + "https://generativelanguage.googleapis.com/v1beta/models/gemini-embedding-2-preview:embedContent" + ); + assert_eq!( + plan.headers.get("x-goog-api-key").map(String::as_str), + Some("sk-upstream-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_eq!(body["content"]["parts"][0]["text"], "hello"); + assert!(body.get("input").is_none()); + assert!(body.get("messages").is_none()); +} + fn embedding_execution_result(plan: &ExecutionPlan) -> ExecutionResult { ExecutionResult { request_id: plan.request_id.clone(), @@ -145,6 +332,55 @@ fn embedding_execution_result(plan: &ExecutionPlan) -> ExecutionResult { } } +fn gemini_embedding_execution_result(plan: &ExecutionPlan) -> ExecutionResult { + ExecutionResult { + request_id: plan.request_id.clone(), + candidate_id: plan.candidate_id.clone(), + status_code: 200, + headers: BTreeMap::from([("content-type".to_string(), "application/json".to_string())]), + body: Some(ResponseBody { + json_body: Some(json!({ + "model": "gemini-embedding-2-preview", + "embedding": { + "values": [0.1, 0.2, 0.3] + }, + "usageMetadata": { + "promptTokenCount": 4, + "totalTokenCount": 4 + } + })), + body_bytes_b64: None, + }), + telemetry: None, + error: None, + } +} + +fn gemini_batch_embedding_execution_result(plan: &ExecutionPlan) -> ExecutionResult { + ExecutionResult { + request_id: plan.request_id.clone(), + candidate_id: plan.candidate_id.clone(), + status_code: 200, + headers: BTreeMap::from([("content-type".to_string(), "application/json".to_string())]), + body: Some(ResponseBody { + json_body: Some(json!({ + "model": "gemini-embedding-2-preview", + "embeddings": [ + {"values": [0.1, 0.2, 0.3]}, + {"values": [0.4, 0.5, 0.6]} + ], + "usageMetadata": { + "promptTokenCount": 8, + "totalTokenCount": 8 + } + })), + body_bytes_b64: None, + }), + telemetry: None, + error: None, + } +} + #[tokio::test] async fn embeddings_route_accepts_openai_payload() { let (execution_runtime_url, execution_runtime_handle) = @@ -215,6 +451,157 @@ async fn embeddings_route_accepts_openai_payload() { execution_runtime_handle.abort(); } +#[tokio::test] +async fn embeddings_route_converts_openai_payload_to_gemini_embedding_provider() { + let (execution_runtime_url, execution_runtime_handle) = + start_server(gemini_embedding_conversion_execution_runtime()).await; + let gateway = build_router_with_state(gemini_embedding_success_state( + execution_runtime_url, + "openai:embedding", + )); + 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-gemini-embedding-success", + ) + .json(&json!({ + "model": "gemini-embedding-2-preview", + "input": "hello" + })) + .send() + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::OK); + assert_eq!( + response + .headers() + .get(CONTROL_ENDPOINT_SIGNATURE_HEADER) + .and_then(|value| value.to_str().ok()), + Some("openai:embedding") + ); + assert_eq!( + response + .headers() + .get(CONTROL_EXECUTION_RUNTIME_HEADER) + .and_then(|value| value.to_str().ok()), + Some("true") + ); + let payload: serde_json::Value = response.json().await.expect("body should parse"); + assert_eq!(payload["object"], "list"); + assert_eq!(payload["model"], "gemini-embedding-2-preview"); + assert_eq!(payload["data"][0]["object"], "embedding"); + assert_eq!(payload["data"][0]["embedding"], json!([0.1, 0.2, 0.3])); + assert_eq!(payload["usage"]["prompt_tokens"], json!(4)); + assert_eq!(payload["usage"]["total_tokens"], json!(4)); + + 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) = + start_server(gemini_embedding_batch_conversion_execution_runtime()).await; + let gateway = build_router_with_state(gemini_embedding_success_state( + execution_runtime_url, + "openai:embedding", + )); + 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-gemini-embedding-success", + ) + .json(&json!({ + "model": "gemini-embedding-2-preview", + "input": ["hello", "world"] + })) + .send() + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::OK); + assert_eq!( + response + .headers() + .get(CONTROL_ENDPOINT_SIGNATURE_HEADER) + .and_then(|value| value.to_str().ok()), + Some("openai:embedding") + ); + let payload: serde_json::Value = response.json().await.expect("body should parse"); + assert_eq!(payload["object"], "list"); + assert_eq!(payload["data"].as_array().map(Vec::len), Some(2)); + assert_eq!(payload["data"][0]["index"], json!(0)); + assert_eq!(payload["data"][0]["embedding"], json!([0.1, 0.2, 0.3])); + assert_eq!(payload["data"][1]["index"], json!(1)); + assert_eq!(payload["data"][1]["embedding"], json!([0.4, 0.5, 0.6])); + assert_eq!(payload["usage"]["prompt_tokens"], json!(8)); + assert_eq!(payload["usage"]["total_tokens"], json!(8)); + + gateway_handle.abort(); + execution_runtime_handle.abort(); +} + +#[tokio::test] +async fn gemini_embed_content_route_uses_native_gemini_embedding_provider() { + let (execution_runtime_url, execution_runtime_handle) = + start_server(gemini_embedding_native_execution_runtime()).await; + let gateway = build_router_with_state(gemini_embedding_success_state( + execution_runtime_url, + "gemini:embedding", + )); + let (gateway_url, gateway_handle) = start_server(gateway).await; + + let response = reqwest::Client::new() + .post(format!( + "{gateway_url}/v1beta/models/gemini-embedding-2-preview:embedContent" + )) + .header("x-goog-api-key", "sk-gemini-embedding-success") + .json(&json!({ + "content": { + "parts": [{"text": "hello"}] + } + })) + .send() + .await + .expect("request should succeed"); + + assert_eq!(response.status(), StatusCode::OK); + assert_eq!( + response + .headers() + .get(CONTROL_ROUTE_FAMILY_HEADER) + .and_then(|value| value.to_str().ok()), + Some("gemini") + ); + assert_eq!( + response + .headers() + .get(CONTROL_ROUTE_KIND_HEADER) + .and_then(|value| value.to_str().ok()), + Some("embedding") + ); + assert_eq!( + response + .headers() + .get(CONTROL_ENDPOINT_SIGNATURE_HEADER) + .and_then(|value| value.to_str().ok()), + Some("gemini:embedding") + ); + let payload: serde_json::Value = response.json().await.expect("body should parse"); + assert_eq!(payload["embedding"]["values"], json!([0.1, 0.2, 0.3])); + assert_eq!(payload["model"], "gemini-embedding-2-preview"); + + gateway_handle.abort(); + execution_runtime_handle.abort(); +} + #[tokio::test] async fn embeddings_route_accepts_all_canonical_input_shapes() { let (execution_runtime_url, execution_runtime_handle) = diff --git a/crates/aether-ai-formats/src/api.rs b/crates/aether-ai-formats/src/api.rs index dc106771a..861369de0 100644 --- a/crates/aether-ai-formats/src/api.rs +++ b/crates/aether-ai-formats/src/api.rs @@ -16,18 +16,21 @@ pub use crate::contracts::{ GEMINI_CHAT_SYNC_SUCCESS_REPORT_KIND, GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_STREAM_SUCCESS_REPORT_KIND, GEMINI_CLI_SYNC_ERROR_REPORT_KIND, GEMINI_CLI_SYNC_FINALIZE_REPORT_KIND, GEMINI_CLI_SYNC_PLAN_KIND, - GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND, GEMINI_FILES_DELETE_PLAN_KIND, + GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, GEMINI_FILES_DELETE_PLAN_KIND, GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, GEMINI_VIDEO_CREATE_SYNC_FINALIZE_REPORT_KIND, GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_STREAM_SUCCESS_REPORT_KIND, OPENAI_CHAT_SYNC_ERROR_REPORT_KIND, OPENAI_CHAT_SYNC_FINALIZE_REPORT_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_CHAT_SYNC_SUCCESS_REPORT_KIND, - OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, - OPENAI_IMAGE_STREAM_SUCCESS_REPORT_KIND, OPENAI_IMAGE_SYNC_ERROR_REPORT_KIND, - OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, - OPENAI_IMAGE_SYNC_SUCCESS_REPORT_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, - OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_SUCCESS_REPORT_KIND, + OPENAI_EMBEDDING_SYNC_ERROR_REPORT_KIND, OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND, + OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, + OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_STREAM_SUCCESS_REPORT_KIND, + OPENAI_IMAGE_SYNC_ERROR_REPORT_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND, + OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_IMAGE_SYNC_SUCCESS_REPORT_KIND, + OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, + OPENAI_RESPONSES_COMPACT_STREAM_SUCCESS_REPORT_KIND, OPENAI_RESPONSES_COMPACT_SYNC_ERROR_REPORT_KIND, OPENAI_RESPONSES_COMPACT_SYNC_FINALIZE_REPORT_KIND, OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_SYNC_SUCCESS_REPORT_KIND, OPENAI_RESPONSES_STREAM_PLAN_KIND, @@ -139,18 +142,22 @@ pub use crate::formats::{ resolve_stream_spec as resolve_gemini_stream_spec, resolve_sync_spec as resolve_gemini_sync_spec, }, - openai::responses::{ - codex::{ - apply_codex_openai_responses_chat_body_edits, - apply_codex_openai_responses_special_body_edits, - apply_codex_openai_responses_special_headers, - apply_openai_responses_compact_special_body_edits, CODEX_OPENAI_IMAGE_DEFAULT_MODEL, - CODEX_OPENAI_IMAGE_DEFAULT_OUTPUT_FORMAT, CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_MODEL, - CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_PROMPT, CODEX_OPENAI_IMAGE_INTERNAL_MODEL, - }, - spec::{ - resolve_stream_spec as resolve_openai_responses_stream_spec, - resolve_sync_spec as resolve_openai_responses_sync_spec, LocalOpenAiResponsesSpec, + openai::{ + embedding::spec::resolve_sync_spec as resolve_openai_embedding_sync_spec, + responses::{ + codex::{ + apply_codex_openai_responses_chat_body_edits, + apply_codex_openai_responses_special_body_edits, + apply_codex_openai_responses_special_headers, + apply_openai_responses_compact_special_body_edits, + CODEX_OPENAI_IMAGE_DEFAULT_MODEL, CODEX_OPENAI_IMAGE_DEFAULT_OUTPUT_FORMAT, + CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_MODEL, + CODEX_OPENAI_IMAGE_DEFAULT_VARIATION_PROMPT, CODEX_OPENAI_IMAGE_INTERNAL_MODEL, + }, + spec::{ + resolve_stream_spec as resolve_openai_responses_stream_spec, + resolve_sync_spec as resolve_openai_responses_sync_spec, LocalOpenAiResponsesSpec, + }, }, }, shared::{ diff --git a/crates/aether-ai-formats/src/contracts/mod.rs b/crates/aether-ai-formats/src/contracts/mod.rs index 769044120..7580a915e 100644 --- a/crates/aether-ai-formats/src/contracts/mod.rs +++ b/crates/aether-ai-formats/src/contracts/mod.rs @@ -14,9 +14,9 @@ pub use plan_kinds::{ is_openai_responses_stream_plan_kind, is_openai_responses_sync_plan_kind, CLAUDE_CHAT_STREAM_PLAN_KIND, CLAUDE_CHAT_SYNC_PLAN_KIND, CLAUDE_CLI_STREAM_PLAN_KIND, CLAUDE_CLI_SYNC_PLAN_KIND, GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, - GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_FILES_DELETE_PLAN_KIND, - GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, - GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, + GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_FILES_DELETE_PLAN_KIND, GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, + GEMINI_FILES_LIST_PLAN_KIND, GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, @@ -36,9 +36,11 @@ pub use report_kinds::{ GEMINI_CHAT_SYNC_ERROR_REPORT_KIND, GEMINI_CHAT_SYNC_FINALIZE_REPORT_KIND, GEMINI_CHAT_SYNC_SUCCESS_REPORT_KIND, GEMINI_CLI_STREAM_SUCCESS_REPORT_KIND, GEMINI_CLI_SYNC_ERROR_REPORT_KIND, GEMINI_CLI_SYNC_FINALIZE_REPORT_KIND, - GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND, GEMINI_VIDEO_CREATE_SYNC_FINALIZE_REPORT_KIND, - OPENAI_CHAT_STREAM_SUCCESS_REPORT_KIND, OPENAI_CHAT_SYNC_ERROR_REPORT_KIND, - OPENAI_CHAT_SYNC_FINALIZE_REPORT_KIND, OPENAI_CHAT_SYNC_SUCCESS_REPORT_KIND, + GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND, GEMINI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, + GEMINI_VIDEO_CREATE_SYNC_FINALIZE_REPORT_KIND, OPENAI_CHAT_STREAM_SUCCESS_REPORT_KIND, + OPENAI_CHAT_SYNC_ERROR_REPORT_KIND, OPENAI_CHAT_SYNC_FINALIZE_REPORT_KIND, + OPENAI_CHAT_SYNC_SUCCESS_REPORT_KIND, OPENAI_EMBEDDING_SYNC_ERROR_REPORT_KIND, + OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND, OPENAI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, OPENAI_IMAGE_STREAM_SUCCESS_REPORT_KIND, OPENAI_IMAGE_SYNC_ERROR_REPORT_KIND, OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND, OPENAI_IMAGE_SYNC_SUCCESS_REPORT_KIND, OPENAI_RESPONSES_COMPACT_STREAM_SUCCESS_REPORT_KIND, diff --git a/crates/aether-ai-formats/src/contracts/plan_kinds.rs b/crates/aether-ai-formats/src/contracts/plan_kinds.rs index afd901ba9..62a24bfb3 100644 --- a/crates/aether-ai-formats/src/contracts/plan_kinds.rs +++ b/crates/aether-ai-formats/src/contracts/plan_kinds.rs @@ -22,6 +22,7 @@ pub const OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND: &str = "openai_video_create_sync"; pub const OPENAI_CHAT_SYNC_PLAN_KIND: &str = "openai_chat_sync"; pub const OPENAI_EMBEDDING_SYNC_PLAN_KIND: &str = "openai_embedding_sync"; pub const OPENAI_RERANK_SYNC_PLAN_KIND: &str = "openai_rerank_sync"; +pub const GEMINI_EMBEDDING_SYNC_PLAN_KIND: &str = "gemini_embedding_sync"; pub const OPENAI_RESPONSES_SYNC_PLAN_KIND: &str = "openai_responses_sync"; pub const OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND: &str = "openai_responses_compact_sync"; pub const CLAUDE_CHAT_SYNC_PLAN_KIND: &str = "claude_chat_sync"; diff --git a/crates/aether-ai-formats/src/contracts/report_kinds.rs b/crates/aether-ai-formats/src/contracts/report_kinds.rs index bf94c57ad..b1e3896ba 100644 --- a/crates/aether-ai-formats/src/contracts/report_kinds.rs +++ b/crates/aether-ai-formats/src/contracts/report_kinds.rs @@ -1,8 +1,8 @@ use crate::contracts::{ CLAUDE_CHAT_SYNC_PLAN_KIND, CLAUDE_CLI_SYNC_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, - GEMINI_CLI_SYNC_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, - OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND, - OPENAI_RESPONSES_SYNC_PLAN_KIND, + GEMINI_CLI_SYNC_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, + OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, + OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND, OPENAI_RESPONSES_SYNC_PLAN_KIND, }; pub const OPENAI_CHAT_SYNC_FINALIZE_REPORT_KIND: &str = "openai_chat_sync_finalize"; @@ -11,6 +11,7 @@ pub const GEMINI_CHAT_SYNC_FINALIZE_REPORT_KIND: &str = "gemini_chat_sync_finali pub const OPENAI_RESPONSES_SYNC_FINALIZE_REPORT_KIND: &str = "openai_responses_sync_finalize"; pub const OPENAI_RESPONSES_COMPACT_SYNC_FINALIZE_REPORT_KIND: &str = "openai_responses_compact_sync_finalize"; +pub const OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND: &str = "openai_embedding_sync_finalize"; pub const OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND: &str = "openai_image_sync_finalize"; pub const CLAUDE_CLI_SYNC_FINALIZE_REPORT_KIND: &str = "claude_cli_sync_finalize"; pub const GEMINI_CLI_SYNC_FINALIZE_REPORT_KIND: &str = "gemini_cli_sync_finalize"; @@ -25,6 +26,8 @@ pub const GEMINI_CHAT_SYNC_SUCCESS_REPORT_KIND: &str = "gemini_chat_sync_success pub const OPENAI_RESPONSES_SYNC_SUCCESS_REPORT_KIND: &str = "openai_responses_sync_success"; pub const OPENAI_RESPONSES_COMPACT_SYNC_SUCCESS_REPORT_KIND: &str = "openai_responses_compact_sync_success"; +pub const OPENAI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND: &str = "openai_embedding_sync_success"; +pub const GEMINI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND: &str = "gemini_embedding_sync_success"; pub const OPENAI_IMAGE_SYNC_SUCCESS_REPORT_KIND: &str = "openai_image_sync_success"; pub const CLAUDE_CLI_SYNC_SUCCESS_REPORT_KIND: &str = "claude_cli_sync_success"; pub const GEMINI_CLI_SYNC_SUCCESS_REPORT_KIND: &str = "gemini_cli_sync_success"; @@ -45,6 +48,7 @@ pub const GEMINI_CHAT_SYNC_ERROR_REPORT_KIND: &str = "gemini_chat_sync_error"; pub const OPENAI_RESPONSES_SYNC_ERROR_REPORT_KIND: &str = "openai_responses_sync_error"; pub const OPENAI_RESPONSES_COMPACT_SYNC_ERROR_REPORT_KIND: &str = "openai_responses_compact_sync_error"; +pub const OPENAI_EMBEDDING_SYNC_ERROR_REPORT_KIND: &str = "openai_embedding_sync_error"; pub const OPENAI_IMAGE_SYNC_ERROR_REPORT_KIND: &str = "openai_image_sync_error"; pub const CLAUDE_CLI_SYNC_ERROR_REPORT_KIND: &str = "claude_cli_sync_error"; pub const GEMINI_CLI_SYNC_ERROR_REPORT_KIND: &str = "gemini_cli_sync_error"; @@ -58,6 +62,7 @@ pub fn implicit_sync_finalize_report_kind(plan_kind: &str) -> Option<&'static st OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND => { Some(OPENAI_RESPONSES_COMPACT_SYNC_FINALIZE_REPORT_KIND) } + OPENAI_EMBEDDING_SYNC_PLAN_KIND => Some(OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND), OPENAI_IMAGE_SYNC_PLAN_KIND => Some(OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND), CLAUDE_CLI_SYNC_PLAN_KIND => Some(CLAUDE_CLI_SYNC_FINALIZE_REPORT_KIND), GEMINI_CLI_SYNC_PLAN_KIND => Some(GEMINI_CLI_SYNC_FINALIZE_REPORT_KIND), @@ -72,6 +77,7 @@ pub fn core_error_default_client_api_format(report_kind: &str) -> Option<&'stati GEMINI_CHAT_SYNC_FINALIZE_REPORT_KIND => Some("gemini:generate_content"), OPENAI_RESPONSES_SYNC_FINALIZE_REPORT_KIND => Some("openai:responses"), OPENAI_RESPONSES_COMPACT_SYNC_FINALIZE_REPORT_KIND => Some("openai:responses:compact"), + OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND => Some("openai:embedding"), LEGACY_OPENAI_CLI_SYNC_FINALIZE_REPORT_KIND => Some("openai:responses"), LEGACY_OPENAI_COMPACT_SYNC_FINALIZE_REPORT_KIND => Some("openai:responses:compact"), OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND => Some("openai:image"), @@ -90,6 +96,7 @@ pub fn core_error_background_report_kind(report_kind: &str) -> Option<&'static s OPENAI_RESPONSES_COMPACT_SYNC_FINALIZE_REPORT_KIND => { Some(OPENAI_RESPONSES_COMPACT_SYNC_ERROR_REPORT_KIND) } + OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND => Some(OPENAI_EMBEDDING_SYNC_ERROR_REPORT_KIND), LEGACY_OPENAI_CLI_SYNC_FINALIZE_REPORT_KIND => { Some(OPENAI_RESPONSES_SYNC_ERROR_REPORT_KIND) } @@ -115,6 +122,9 @@ pub fn core_success_background_report_kind(report_kind: &str) -> Option<&'static OPENAI_RESPONSES_COMPACT_SYNC_FINALIZE_REPORT_KIND => { Some(OPENAI_RESPONSES_COMPACT_SYNC_SUCCESS_REPORT_KIND) } + OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND => { + Some(OPENAI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND) + } LEGACY_OPENAI_CLI_SYNC_FINALIZE_REPORT_KIND => { Some(OPENAI_RESPONSES_SYNC_SUCCESS_REPORT_KIND) } diff --git a/crates/aether-ai-formats/src/formats/gemini/embedding/mod.rs b/crates/aether-ai-formats/src/formats/gemini/embedding/mod.rs index be9378d97..e0062185c 100644 --- a/crates/aether-ai-formats/src/formats/gemini/embedding/mod.rs +++ b/crates/aether-ai-formats/src/formats/gemini/embedding/mod.rs @@ -1 +1,2 @@ pub mod request; +pub mod response; diff --git a/crates/aether-ai-formats/src/formats/gemini/embedding/request.rs b/crates/aether-ai-formats/src/formats/gemini/embedding/request.rs index 7b569fb2a..83416a830 100644 --- a/crates/aether-ai-formats/src/formats/gemini/embedding/request.rs +++ b/crates/aether-ai-formats/src/formats/gemini/embedding/request.rs @@ -1,9 +1,8 @@ -use serde_json::json; -use serde_json::Value; +use serde_json::{json, Map, Value}; use crate::formats::context::FormatContext; use crate::formats::openai::embedding::request::mapped_embedding_model; -use crate::protocol::canonical::CanonicalRequest; +use crate::protocol::canonical::{CanonicalEmbeddingRequest, CanonicalRequest}; pub fn to(request: &CanonicalRequest, ctx: &FormatContext) -> Option { let embedding = request.embedding.as_ref()?; @@ -13,22 +12,169 @@ pub fn to(request: &CanonicalRequest, ctx: &FormatContext) -> Option { } let model = mapped_embedding_model(request, ctx.mapped_model_or(request.model.as_str())); if items.len() == 1 { - return Some(json!({ - "model": model, - "content": { - "parts": [{"text": items[0]}] - } - })); + return Some(Value::Object(gemini_embedding_request_object( + &model, items[0], embedding, + ))); + } + let model_resource = gemini_embedding_model_resource_name(&model); + let requests = items + .into_iter() + .map(|text| { + Value::Object(gemini_embedding_request_object( + &model_resource, + text, + embedding, + )) + }) + .collect::>(); + Some(json!({ "requests": requests })) +} + +fn gemini_embedding_request_object( + model: &str, + text: &str, + embedding: &CanonicalEmbeddingRequest, +) -> Map { + let mut object = Map::new(); + object.insert("model".to_string(), Value::String(model.to_string())); + object.insert( + "content".to_string(), + json!({ + "parts": [{"text": text}] + }), + ); + insert_gemini_embedding_options(&mut object, embedding); + object +} + +fn gemini_embedding_model_resource_name(model: &str) -> String { + let trimmed = model.trim(); + if trimmed.starts_with("models/") { + trimmed.to_string() + } else { + format!("models/{trimmed}") + } +} + +fn insert_gemini_embedding_options( + object: &mut Map, + embedding: &CanonicalEmbeddingRequest, +) { + if let Some(dimensions) = embedding.dimensions { + object.insert("outputDimensionality".to_string(), Value::from(dimensions)); + } + if let Some(task_type) = embedding + .task + .as_deref() + .and_then(normalize_gemini_embedding_task_type) + { + object.insert("taskType".to_string(), Value::String(task_type)); + } +} + +fn normalize_gemini_embedding_task_type(value: &str) -> Option { + let normalized = value.trim(); + if normalized.is_empty() { + return None; + } + let key = normalized.replace(['-', ' '], "_").to_ascii_uppercase(); + let task_type = match key.as_str() { + "QUERY" | "RETRIEVAL_QUERY" => "RETRIEVAL_QUERY", + "DOCUMENT" | "RETRIEVAL_DOCUMENT" => "RETRIEVAL_DOCUMENT", + "TEXT_MATCHING" | "SEMANTIC_SIMILARITY" => "SEMANTIC_SIMILARITY", + "CLASSIFICATION" => "CLASSIFICATION", + "CLUSTERING" => "CLUSTERING", + "QUESTION_ANSWERING" => "QUESTION_ANSWERING", + "FACT_VERIFICATION" => "FACT_VERIFICATION", + "CODE_RETRIEVAL_QUERY" => "CODE_RETRIEVAL_QUERY", + _ => key.as_str(), + }; + Some(task_type.to_string()) +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + + use serde_json::json; + + use super::to; + use crate::formats::context::FormatContext; + use crate::protocol::canonical::{ + CanonicalEmbeddingInput, CanonicalEmbeddingRequest, CanonicalRequest, + }; + + fn canonical_embedding(input: CanonicalEmbeddingInput) -> CanonicalRequest { + CanonicalRequest { + model: "text-embedding-3-small".to_string(), + embedding: Some(CanonicalEmbeddingRequest { + input, + encoding_format: None, + dimensions: None, + task: None, + user: None, + extensions: BTreeMap::new(), + }), + ..CanonicalRequest::default() + } + } + + #[test] + fn single_string_array_item_uses_single_embed_content_body() { + let request = canonical_embedding(CanonicalEmbeddingInput::StringArray(vec![ + "hello".to_string() + ])); + let body = to( + &request, + &FormatContext::default().with_mapped_model("gemini-embedding-2-preview"), + ) + .expect("gemini embedding request"); + + assert_eq!(body["model"], "gemini-embedding-2-preview"); + assert_eq!(body["content"]["parts"][0]["text"], "hello"); + assert!(body.get("requests").is_none()); + } + + #[test] + fn multiple_string_items_use_gemini_batch_request_body() { + let request = canonical_embedding(CanonicalEmbeddingInput::StringArray(vec![ + "alpha".to_string(), + "beta".to_string(), + ])); + let body = to( + &request, + &FormatContext::default().with_mapped_model("gemini-embedding-2-preview"), + ) + .expect("gemini embedding request"); + + assert!(body.get("model").is_none()); + assert_eq!(body["requests"].as_array().map(Vec::len), Some(2)); + assert_eq!( + body["requests"][0]["model"], + "models/gemini-embedding-2-preview" + ); + assert_eq!(body["requests"][0]["content"]["parts"][0]["text"], "alpha"); + assert_eq!( + body["requests"][1]["model"], + "models/gemini-embedding-2-preview" + ); + assert_eq!(body["requests"][1]["content"]["parts"][0]["text"], "beta"); + } + + #[test] + fn explicit_embedding_options_are_preserved_without_defaults() { + let mut request = canonical_embedding(CanonicalEmbeddingInput::String("query".to_string())); + let embedding = request.embedding.as_mut().expect("embedding request"); + embedding.dimensions = Some(768); + embedding.task = Some("retrieval_query".to_string()); + + let body = to( + &request, + &FormatContext::default().with_mapped_model("gemini-embedding-2-preview"), + ) + .expect("gemini embedding request"); + + assert_eq!(body["outputDimensionality"], json!(768)); + assert_eq!(body["taskType"], "RETRIEVAL_QUERY"); } - Some(json!({ - "model": model, - "requests": items.into_iter().map(|text| { - json!({ - "model": model, - "content": { - "parts": [{"text": text}] - } - }) - }).collect::>() - })) } diff --git a/crates/aether-ai-formats/src/formats/gemini/embedding/response.rs b/crates/aether-ai-formats/src/formats/gemini/embedding/response.rs new file mode 100644 index 000000000..91ccf5e0c --- /dev/null +++ b/crates/aether-ai-formats/src/formats/gemini/embedding/response.rs @@ -0,0 +1,108 @@ +use serde_json::Value; + +use crate::formats::openai::embedding::request::namespace_extensions; +use crate::protocol::canonical::{ + gemini_usage_to_canonical, CanonicalEmbedding, CanonicalEmbeddingResponse, +}; + +pub fn from(body_json: &Value) -> Option { + let body = body_json.as_object()?; + if body.contains_key("error") { + return None; + } + + let embeddings = if let Some(values) = body + .get("embedding") + .and_then(Value::as_object) + .and_then(|embedding| embedding.get("values")) + .and_then(Value::as_array) + { + vec![CanonicalEmbedding { + index: 0, + embedding: embedding_values(values)?, + extensions: Default::default(), + }] + } else { + let raw_embeddings = body.get("embeddings")?.as_array()?; + raw_embeddings + .iter() + .enumerate() + .map(|(index, item)| { + let item_object = item.as_object()?; + let values = item_object.get("values")?.as_array()?; + Some(CanonicalEmbedding { + index, + embedding: embedding_values(values)?, + extensions: namespace_extensions("gemini", item_object, &["values"]), + }) + }) + .collect::>>()? + }; + if embeddings.is_empty() + || embeddings + .iter() + .any(|embedding| embedding.embedding.is_empty()) + { + return None; + } + + Some(CanonicalEmbeddingResponse { + id: body + .get("id") + .or_else(|| body.get("responseId")) + .and_then(Value::as_str) + .unwrap_or("embd-gemini-unknown") + .to_string(), + model: body + .get("model") + .or_else(|| body.get("modelVersion")) + .and_then(Value::as_str) + .unwrap_or("unknown") + .to_string(), + embeddings, + usage: gemini_usage_to_canonical(body.get("usageMetadata")), + extensions: namespace_extensions( + "gemini", + body, + &[ + "id", + "responseId", + "model", + "modelVersion", + "embedding", + "embeddings", + "usageMetadata", + ], + ), + }) +} + +fn embedding_values(values: &[Value]) -> Option> { + values.iter().map(Value::as_f64).collect() +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::from; + + #[test] + fn parses_gemini_single_embedding_response() { + let body = json!({ + "embedding": {"values": [0.1, 0.2, 0.3]}, + "usageMetadata": { + "promptTokenCount": 4, + "totalTokenCount": 4 + } + }); + + let parsed = from(&body).expect("response should parse"); + + assert_eq!(parsed.model, "unknown"); + assert_eq!(parsed.embeddings[0].embedding, vec![0.1, 0.2, 0.3]); + let usage = parsed.usage.expect("usage should parse"); + assert_eq!(usage.input_tokens, 4); + assert_eq!(usage.total_tokens, 4); + } +} diff --git a/crates/aether-ai-formats/src/formats/openai/embedding/mod.rs b/crates/aether-ai-formats/src/formats/openai/embedding/mod.rs index e0062185c..75cf53e33 100644 --- a/crates/aether-ai-formats/src/formats/openai/embedding/mod.rs +++ b/crates/aether-ai-formats/src/formats/openai/embedding/mod.rs @@ -1,2 +1,3 @@ pub mod request; pub mod response; +pub mod spec; diff --git a/crates/aether-ai-formats/src/formats/openai/embedding/spec.rs b/crates/aether-ai-formats/src/formats/openai/embedding/spec.rs new file mode 100644 index 000000000..b3f0bdc9c --- /dev/null +++ b/crates/aether-ai-formats/src/formats/openai/embedding/spec.rs @@ -0,0 +1,35 @@ +use crate::contracts::{ + OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, +}; +use crate::formats::shared::family::{ + LocalStandardSourceFamily, LocalStandardSourceMode, LocalStandardSpec, +}; + +pub fn resolve_sync_spec(plan_kind: &str) -> Option { + match plan_kind { + OPENAI_EMBEDDING_SYNC_PLAN_KIND => Some(LocalStandardSpec { + api_format: "openai:embedding", + decision_kind: OPENAI_EMBEDDING_SYNC_PLAN_KIND, + report_kind: OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND, + family: LocalStandardSourceFamily::Standard, + mode: LocalStandardSourceMode::Embedding, + require_streaming: false, + }), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::resolve_sync_spec; + use crate::formats::shared::family::LocalStandardSourceMode; + + #[test] + fn resolves_openai_embedding_sync_standard_spec() { + let spec = resolve_sync_spec("openai_embedding_sync").expect("spec"); + assert_eq!(spec.api_format, "openai:embedding"); + assert_eq!(spec.report_kind, "openai_embedding_sync_finalize"); + assert_eq!(spec.mode, LocalStandardSourceMode::Embedding); + assert!(!spec.require_streaming); + } +} diff --git a/crates/aether-ai-formats/src/formats/registry.rs b/crates/aether-ai-formats/src/formats/registry.rs index babc18537..ff06e8c29 100644 --- a/crates/aether-ai-formats/src/formats/registry.rs +++ b/crates/aether-ai-formats/src/formats/registry.rs @@ -228,11 +228,16 @@ mod tests { &FormatContext::default().with_mapped_model("gemini-embedding-001"), ) .expect("gemini embedding conversion should succeed"); - assert_eq!(gemini["model"], "gemini-embedding-001"); + assert!(gemini.get("model").is_none()); + assert_eq!( + gemini["requests"][0]["model"], + "models/gemini-embedding-001" + ); assert_eq!( gemini["requests"][0]["content"]["parts"][0]["text"], "alpha" ); + assert_eq!(gemini["requests"][0]["outputDimensionality"], 2); assert!(gemini.get("messages").is_none()); let doubao = convert_request( diff --git a/crates/aether-ai-formats/src/formats/shared/error_body.rs b/crates/aether-ai-formats/src/formats/shared/error_body.rs index d75e5003a..53ace4d28 100644 --- a/crates/aether-ai-formats/src/formats/shared/error_body.rs +++ b/crates/aether-ai-formats/src/formats/shared/error_body.rs @@ -38,7 +38,7 @@ pub fn build_core_error_body_for_client_format( error_object.insert("message".to_string(), Value::String(message.to_string())); match aether_ai_formats::normalize_api_format_alias(client_api_format).as_str() { - "openai:chat" | "openai:responses" | "openai:responses:compact" => { + "openai:chat" | "openai:responses" | "openai:responses:compact" | "openai:embedding" => { error_object.insert( "type".to_string(), Value::String(map_local_sync_error_kind_to_openai_type(kind).to_string()), diff --git a/crates/aether-ai-formats/src/formats/shared/passthrough.rs b/crates/aether-ai-formats/src/formats/shared/passthrough.rs index e37152fd0..63de6080c 100644 --- a/crates/aether-ai-formats/src/formats/shared/passthrough.rs +++ b/crates/aether-ai-formats/src/formats/shared/passthrough.rs @@ -1,7 +1,8 @@ use crate::contracts::{ CLAUDE_CHAT_STREAM_PLAN_KIND, CLAUDE_CHAT_SYNC_PLAN_KIND, CLAUDE_CLI_STREAM_PLAN_KIND, CLAUDE_CLI_SYNC_PLAN_KIND, GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, - GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, }; @@ -50,6 +51,13 @@ pub fn resolve_sync_spec(plan_kind: &str) -> Option family: LocalSameFormatProviderFamily::Gemini, require_streaming: false, }), + GEMINI_EMBEDDING_SYNC_PLAN_KIND => Some(LocalSameFormatProviderSpec { + api_format: "gemini:embedding", + decision_kind: GEMINI_EMBEDDING_SYNC_PLAN_KIND, + report_kind: GEMINI_EMBEDDING_SYNC_SUCCESS_REPORT_KIND, + family: LocalSameFormatProviderFamily::Gemini, + require_streaming: false, + }), OPENAI_EMBEDDING_SYNC_PLAN_KIND => Some(LocalSameFormatProviderSpec { api_format: "openai:embedding", decision_kind: OPENAI_EMBEDDING_SYNC_PLAN_KIND, @@ -130,6 +138,15 @@ mod tests { assert!(!spec.require_streaming); } + #[test] + fn resolves_gemini_embedding_sync_same_format_spec() { + let spec = resolve_sync_spec("gemini_embedding_sync").expect("spec"); + assert_eq!(spec.api_format, "gemini:embedding"); + assert_eq!(spec.report_kind, "gemini_embedding_sync_success"); + assert_eq!(spec.family, super::LocalSameFormatProviderFamily::Gemini); + assert!(!spec.require_streaming); + } + #[test] fn resolves_openai_rerank_sync_same_format_spec() { let spec = resolve_sync_spec("openai_rerank_sync").expect("spec"); diff --git a/crates/aether-ai-formats/src/formats/shared/routing.rs b/crates/aether-ai-formats/src/formats/shared/routing.rs index 0711881ee..dbec2843f 100644 --- a/crates/aether-ai-formats/src/formats/shared/routing.rs +++ b/crates/aether-ai-formats/src/formats/shared/routing.rs @@ -4,9 +4,9 @@ use url::form_urlencoded; use crate::contracts::{ CLAUDE_CHAT_STREAM_PLAN_KIND, CLAUDE_CHAT_SYNC_PLAN_KIND, CLAUDE_CLI_STREAM_PLAN_KIND, CLAUDE_CLI_SYNC_PLAN_KIND, GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, - GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_FILES_DELETE_PLAN_KIND, - GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, GEMINI_FILES_LIST_PLAN_KIND, - GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, + GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + GEMINI_FILES_DELETE_PLAN_KIND, GEMINI_FILES_DOWNLOAD_PLAN_KIND, GEMINI_FILES_GET_PLAN_KIND, + GEMINI_FILES_LIST_PLAN_KIND, GEMINI_FILES_UPLOAD_PLAN_KIND, GEMINI_VIDEO_CANCEL_SYNC_PLAN_KIND, GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, @@ -166,6 +166,14 @@ pub fn resolve_execution_runtime_sync_plan_kind( return Some(GEMINI_VIDEO_CREATE_SYNC_PLAN_KIND); } + if route_family == Some("gemini") + && route_kind == Some("embedding") + && *method == Method::POST + && (path.ends_with(":embedContent") || path.ends_with(":batchEmbedContents")) + { + return Some(GEMINI_EMBEDDING_SYNC_PLAN_KIND); + } + if route_family == Some("openai") && route_kind == Some("chat") && *method == Method::POST @@ -414,6 +422,7 @@ pub fn supports_sync_execution_decision_kind(plan_kind: &str) -> bool { | CLAUDE_CLI_SYNC_PLAN_KIND | GEMINI_CHAT_SYNC_PLAN_KIND | GEMINI_CLI_SYNC_PLAN_KIND + | GEMINI_EMBEDDING_SYNC_PLAN_KIND | GEMINI_FILES_UPLOAD_PLAN_KIND | OPENAI_VIDEO_CREATE_SYNC_PLAN_KIND | OPENAI_VIDEO_REMIX_SYNC_PLAN_KIND @@ -458,9 +467,9 @@ mod tests { use crate::contracts::{ CLAUDE_CHAT_STREAM_PLAN_KIND, CLAUDE_CHAT_SYNC_PLAN_KIND, CLAUDE_CLI_STREAM_PLAN_KIND, CLAUDE_CLI_SYNC_PLAN_KIND, GEMINI_CHAT_STREAM_PLAN_KIND, GEMINI_CHAT_SYNC_PLAN_KIND, - GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, OPENAI_CHAT_STREAM_PLAN_KIND, - OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, OPENAI_IMAGE_STREAM_PLAN_KIND, - OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, + GEMINI_CLI_STREAM_PLAN_KIND, GEMINI_CLI_SYNC_PLAN_KIND, GEMINI_EMBEDDING_SYNC_PLAN_KIND, + OPENAI_CHAT_STREAM_PLAN_KIND, OPENAI_CHAT_SYNC_PLAN_KIND, OPENAI_EMBEDDING_SYNC_PLAN_KIND, + OPENAI_IMAGE_STREAM_PLAN_KIND, OPENAI_IMAGE_SYNC_PLAN_KIND, OPENAI_RERANK_SYNC_PLAN_KIND, OPENAI_RESPONSES_COMPACT_STREAM_PLAN_KIND, OPENAI_RESPONSES_COMPACT_SYNC_PLAN_KIND, OPENAI_RESPONSES_STREAM_PLAN_KIND, OPENAI_RESPONSES_SYNC_PLAN_KIND, }; @@ -786,6 +795,35 @@ mod tests { )); } + #[test] + fn resolves_gemini_embedding_sync_plan_kind() { + assert_eq!( + resolve_execution_runtime_sync_plan_kind( + Some("ai_public"), + Some("gemini"), + Some("embedding"), + Some("api_key"), + &Method::POST, + "/v1beta/models/gemini-embedding-2-preview:embedContent", + ), + Some(GEMINI_EMBEDDING_SYNC_PLAN_KIND) + ); + assert_eq!( + resolve_execution_runtime_sync_plan_kind( + Some("ai_public"), + Some("gemini"), + Some("embedding"), + Some("api_key"), + &Method::POST, + "/v1beta/models/gemini-embedding-2-preview:batchEmbedContents", + ), + Some(GEMINI_EMBEDDING_SYNC_PLAN_KIND) + ); + assert!(supports_sync_execution_decision_kind( + GEMINI_EMBEDDING_SYNC_PLAN_KIND + )); + } + #[test] fn resolves_openai_rerank_sync_plan_kind() { assert_eq!( diff --git a/crates/aether-ai-formats/src/formats/shared/sync_products.rs b/crates/aether-ai-formats/src/formats/shared/sync_products.rs index db3349ee7..86814f63d 100644 --- a/crates/aether-ai-formats/src/formats/shared/sync_products.rs +++ b/crates/aether-ai-formats/src/formats/shared/sync_products.rs @@ -9,9 +9,10 @@ use aether_ai_formats::formats::conversion::response::{ }; use aether_ai_formats::formats::registry::{convert_response, FormatContext}; use aether_ai_formats::{ - canonical_to_claude_response, canonical_to_gemini_response, canonical_to_openai_chat_response, - canonical_to_openai_responses_compact_response, canonical_to_openai_responses_response, - from_claude_to_canonical_response, from_gemini_to_canonical_response, + canonical_to_claude_response, canonical_to_embedding_response, canonical_to_gemini_response, + canonical_to_openai_chat_response, canonical_to_openai_responses_compact_response, + canonical_to_openai_responses_response, from_claude_to_canonical_response, + from_embedding_to_canonical_response, from_gemini_to_canonical_response, from_openai_chat_to_canonical_response, from_openai_responses_to_canonical_response, sync_chat_response_conversion_kind, sync_cli_response_conversion_kind, }; @@ -330,6 +331,18 @@ pub fn maybe_build_standard_sync_finalize_product_from_normalized_payload( ))); } + if let Some(product) = maybe_build_embedding_cross_format_sync_product_from_normalized_payload( + report_kind, + status_code, + report_context, + body_json, + body_base64, + )? { + return Ok(Some(StandardSyncFinalizeNormalizedProduct::CrossFormat( + product, + ))); + } + Ok( maybe_build_standard_cross_format_sync_product_from_normalized_payload( report_kind, @@ -342,6 +355,90 @@ pub fn maybe_build_standard_sync_finalize_product_from_normalized_payload( ) } +pub fn maybe_build_embedding_cross_format_sync_product_from_normalized_payload( + report_kind: &str, + status_code: u16, + report_context: Option<&Value>, + body_json: Option<&Value>, + body_base64: Option<&str>, +) -> Result, AiSurfaceFinalizeError> { + if report_kind != crate::contracts::OPENAI_EMBEDDING_SYNC_FINALIZE_REPORT_KIND + || status_code >= 400 + { + return Ok(None); + } + + let Some(report_context) = report_context else { + return Ok(None); + }; + let provider_api_format = report_context + .get("provider_api_format") + .and_then(Value::as_str) + .unwrap_or_default() + .trim() + .to_ascii_lowercase(); + let client_api_format = report_context + .get("client_api_format") + .and_then(Value::as_str) + .unwrap_or_default() + .trim() + .to_ascii_lowercase(); + + if provider_api_format == client_api_format { + return Ok(None); + } + + let Some(provider_namespace) = + embedding_response_namespace_for_api_format(&provider_api_format) + else { + return Ok(None); + }; + let Some(client_namespace) = embedding_response_namespace_for_api_format(&client_api_format) + else { + return Ok(None); + }; + + let provider_body_json = match body_base64 { + Some(body_base64) => { + let body_bytes = base64::engine::general_purpose::STANDARD.decode(body_base64)?; + serde_json::from_slice::(&body_bytes).ok() + } + None => body_json.cloned(), + }; + let Some(provider_body_json) = + provider_body_json.filter(|value| !is_error_like_sync_body(value)) + else { + return Ok(None); + }; + + let mut canonical = + match from_embedding_to_canonical_response(&provider_body_json, provider_namespace) { + Some(canonical) => canonical, + None => return Ok(None), + }; + apply_report_context_model_fallback(&mut canonical.model, report_context); + let Some(client_body_json) = canonical_to_embedding_response(&canonical, client_namespace) + else { + return Ok(None); + }; + let client_body_json = + client_body_with_report_context_model(client_body_json, report_context, &client_api_format); + + Ok(Some(StandardCrossFormatSyncProduct { + client_body_json, + provider_body_json, + })) +} + +fn embedding_response_namespace_for_api_format(api_format: &str) -> Option<&'static str> { + match aether_ai_formats::normalize_api_format_alias(api_format).as_str() { + "openai:embedding" => Some("openai"), + "jina:embedding" => Some("jina"), + "gemini:embedding" => Some("gemini"), + _ => None, + } +} + fn maybe_build_standard_same_format_sync_body( report_kind: &str, status_code: u16, @@ -829,6 +926,9 @@ fn client_body_with_report_context_model( "openai:chat" | "openai:responses" | "openai:responses:compact" | "claude:messages" => { object.insert("model".to_string(), Value::String(display_model)); } + "openai:embedding" | "jina:embedding" => { + object.insert("model".to_string(), Value::String(display_model)); + } "gemini:generate_content" => { object.insert("modelVersion".to_string(), Value::String(display_model)); } @@ -1310,6 +1410,7 @@ fn standard_same_format_api_format(report_kind: &str) -> Option<&'static str> { "openai_chat_sync_finalize" => Some("openai:chat"), "claude_chat_sync_finalize" => Some("claude:messages"), "gemini_chat_sync_finalize" => Some("gemini:generate_content"), + "openai_embedding_sync_finalize" => Some("openai:embedding"), "claude_cli_sync_finalize" => Some("claude:messages"), "gemini_cli_sync_finalize" => Some("gemini:generate_content"), _ => None, diff --git a/crates/aether-ai-formats/src/protocol/canonical.rs b/crates/aether-ai-formats/src/protocol/canonical.rs index dc903b9bd..cde609407 100644 --- a/crates/aether-ai-formats/src/protocol/canonical.rs +++ b/crates/aether-ai-formats/src/protocol/canonical.rs @@ -684,6 +684,7 @@ pub fn from_embedding_to_canonical_response( crate::formats::openai::embedding::response::from_namespace(body_json, "openai") } "jina" => crate::formats::openai::embedding::response::from_namespace(body_json, "jina"), + "gemini" => crate::formats::gemini::embedding::response::from(body_json), _ => None, } } @@ -4752,10 +4753,16 @@ mod tests { let gemini = super::canonical_to_embedding_request(&canonical, "gemini-embedding-001", "gemini") .expect("gemini embedding request"); + assert!(gemini.get("model").is_none()); + assert_eq!( + gemini["requests"][0]["model"], + "models/gemini-embedding-001" + ); assert_eq!( gemini["requests"][0]["content"]["parts"][0]["text"], "alpha" ); + assert_eq!(gemini["requests"][0]["outputDimensionality"], 2); assert!(gemini.get("messages").is_none()); let doubao = super::canonical_to_embedding_request( diff --git a/crates/aether-provider-transport/src/conversion.rs b/crates/aether-provider-transport/src/conversion.rs index ca2eee161..059606887 100644 --- a/crates/aether-provider-transport/src/conversion.rs +++ b/crates/aether-provider-transport/src/conversion.rs @@ -1,5 +1,6 @@ use aether_ai_formats::formats::matrix::{ - request_conversion_kind, request_conversion_requires_enable_flag, RequestConversionKind, + api_data_format_id, request_conversion_kind, request_conversion_requires_enable_flag, + RequestConversionKind, }; use aether_ai_formats::normalize_api_format_alias; @@ -39,7 +40,14 @@ pub fn request_conversion_enabled_for_transport( if client_api_format == provider_api_format { return true; } - if request_conversion_kind(client_api_format.as_str(), provider_api_format.as_str()).is_none() { + let conversion_kind = + request_conversion_kind(client_api_format.as_str(), provider_api_format.as_str()); + if conversion_kind.is_none() + && !same_data_format_transport_pair( + client_api_format.as_str(), + provider_api_format.as_str(), + ) + { return false; } if !request_conversion_requires_enable_flag( @@ -62,7 +70,14 @@ pub fn request_pair_allowed_for_transport( if client_api_format == provider_api_format { return true; } - if request_conversion_kind(client_api_format.as_str(), provider_api_format.as_str()).is_none() { + let conversion_kind = + request_conversion_kind(client_api_format.as_str(), provider_api_format.as_str()); + if conversion_kind.is_none() + && !same_data_format_transport_pair( + client_api_format.as_str(), + provider_api_format.as_str(), + ) + { return false; } if is_kiro_claude_messages_transport(transport, &provider_api_format) { @@ -80,6 +95,19 @@ pub fn request_pair_allowed_for_transport( ) } +fn same_data_format_transport_pair(client_api_format: &str, provider_api_format: &str) -> bool { + if aether_ai_formats::api_format_alias_matches(client_api_format, provider_api_format) { + return false; + } + matches!( + ( + api_data_format_id(client_api_format), + api_data_format_id(provider_api_format) + ), + (Some("embedding"), Some("embedding")) | (Some("rerank"), Some("rerank")) + ) +} + pub fn request_conversion_transport_supported( transport: &GatewayProviderTransportSnapshot, kind: RequestConversionKind, @@ -117,15 +145,72 @@ pub fn request_conversion_transport_unsupported_reason( } } +pub fn request_pair_transport_unsupported_reason( + transport: &GatewayProviderTransportSnapshot, + client_api_format: &str, + provider_api_format: &str, +) -> Option<&'static str> { + let client_api_format = normalize_api_format_alias(client_api_format); + let provider_api_format = normalize_api_format_alias(provider_api_format); + + if let Some(kind) = + request_conversion_kind(client_api_format.as_str(), provider_api_format.as_str()) + { + return request_conversion_transport_unsupported_reason(transport, kind); + } + + if !same_data_format_transport_pair(client_api_format.as_str(), provider_api_format.as_str()) { + return Some("transport_api_format_unsupported"); + } + + match provider_api_format.as_str() { + "gemini:embedding" => { + if is_vertex_transport_context(transport) { + local_vertex_gemini_transport_unsupported_reason_with_network(transport) + } else { + local_gemini_transport_unsupported_reason_with_network( + transport, + "gemini:embedding", + ) + } + } + "openai:embedding" | "jina:embedding" | "doubao:embedding" | "openai:rerank" + | "jina:rerank" => local_standard_transport_unsupported_reason_with_network( + transport, + provider_api_format.as_str(), + ), + _ => Some("transport_api_format_unsupported"), + } +} + pub fn request_conversion_direct_auth( transport: &GatewayProviderTransportSnapshot, _kind: RequestConversionKind, ) -> Option<(String, String)> { - match normalize_api_format_alias(&transport.endpoint.api_format).as_str() { - "openai:chat" | "openai:responses" | "openai:responses:compact" => { - resolve_local_openai_bearer_auth(transport) - } - "gemini:generate_content" => { + request_direct_auth_for_provider_format(transport, transport.endpoint.api_format.as_str()) +} + +pub fn request_pair_direct_auth( + transport: &GatewayProviderTransportSnapshot, + provider_api_format: &str, +) -> Option<(String, String)> { + request_direct_auth_for_provider_format(transport, provider_api_format) +} + +fn request_direct_auth_for_provider_format( + transport: &GatewayProviderTransportSnapshot, + provider_api_format: &str, +) -> Option<(String, String)> { + match normalize_api_format_alias(provider_api_format).as_str() { + "openai:chat" + | "openai:responses" + | "openai:responses:compact" + | "openai:embedding" + | "jina:embedding" + | "doubao:embedding" + | "openai:rerank" + | "jina:rerank" => resolve_local_openai_bearer_auth(transport), + "gemini:generate_content" | "gemini:embedding" => { if is_vertex_api_key_transport_context(transport) { resolve_local_vertex_api_key_query_auth(transport) .map(|auth| (VERTEX_API_KEY_QUERY_PARAM.to_string(), auth.value)) diff --git a/crates/aether-provider-transport/src/lib.rs b/crates/aether-provider-transport/src/lib.rs index d7bd27e58..d00471001 100644 --- a/crates/aether-provider-transport/src/lib.rs +++ b/crates/aether-provider-transport/src/lib.rs @@ -30,7 +30,8 @@ pub use conversion::{ candidate_common_transport_skip_reason, candidate_transport_pair_skip_reason, request_conversion_direct_auth, request_conversion_enabled_for_transport, request_conversion_transport_supported, request_conversion_transport_unsupported_reason, - request_pair_allowed_for_transport, CandidateTransportPolicyFacts, + request_pair_allowed_for_transport, request_pair_direct_auth, + request_pair_transport_unsupported_reason, CandidateTransportPolicyFacts, }; pub use diagnostics::{ append_transport_diagnostics_to_value, build_request_trace_proxy_value, @@ -72,6 +73,7 @@ pub use request_url::{ build_cross_format_openai_chat_upstream_url, build_cross_format_openai_responses_upstream_url, build_kiro_cross_format_upstream_url, build_local_openai_chat_upstream_url, build_local_openai_responses_upstream_url, build_transport_request_url, + build_transport_request_url_for_request_body, gemini_embedding_request_body_uses_batch, TransportRequestUrlParams, }; pub use rules::{ diff --git a/crates/aether-provider-transport/src/request_url/mod.rs b/crates/aether-provider-transport/src/request_url/mod.rs index 6a8c8642f..842527963 100644 --- a/crates/aether-provider-transport/src/request_url/mod.rs +++ b/crates/aether-provider-transport/src/request_url/mod.rs @@ -2,6 +2,7 @@ use std::collections::BTreeMap; use std::sync::OnceLock; use regex::Regex; +use serde_json::Value; use url::form_urlencoded; use crate::antigravity::{ @@ -31,6 +32,35 @@ pub struct TransportRequestUrlParams<'a> { pub fn build_transport_request_url( transport: &GatewayProviderTransportSnapshot, params: TransportRequestUrlParams<'_>, +) -> Option { + build_transport_request_url_inner(transport, params, false) +} + +pub fn build_transport_request_url_for_request_body( + transport: &GatewayProviderTransportSnapshot, + params: TransportRequestUrlParams<'_>, + provider_request_body: Option<&Value>, +) -> Option { + let gemini_embedding_batch = + gemini_embedding_request_body_uses_batch(params.provider_api_format, provider_request_body); + build_transport_request_url_inner(transport, params, gemini_embedding_batch) +} + +pub fn gemini_embedding_request_body_uses_batch( + provider_api_format: &str, + provider_request_body: Option<&Value>, +) -> bool { + aether_ai_formats::normalize_api_format_alias(provider_api_format) == "gemini:embedding" + && provider_request_body + .and_then(|body| body.get("requests")) + .and_then(Value::as_array) + .is_some_and(|requests| !requests.is_empty()) +} + +fn build_transport_request_url_inner( + transport: &GatewayProviderTransportSnapshot, + params: TransportRequestUrlParams<'_>, + gemini_embedding_batch: bool, ) -> Option { if let Some(url) = build_transport_hook_url(transport, params) { return Some(url); @@ -45,7 +75,9 @@ pub fn build_transport_request_url( .as_deref() .map(str::trim) .filter(|value| !value.is_empty()) - .map(|path| expand_custom_path_template(path, build_path_params(params))); + .map(|path| { + expand_custom_path_template(path, build_path_params(params, gemini_embedding_batch)) + }); if let Some(path) = custom_path.as_deref() { let blocked_keys = if normalized_provider_api_format.starts_with("gemini:") { @@ -55,6 +87,8 @@ pub fn build_transport_request_url( }; let normalized_path = if normalized_provider_api_format == "gemini:generate_content" { normalize_gemini_content_action_path(path, params.upstream_is_stream) + } else if normalized_provider_api_format == "gemini:embedding" { + normalize_gemini_embedding_action_path(path, gemini_embedding_batch) } else { path.to_string() }; @@ -106,6 +140,7 @@ pub fn build_transport_request_url( &transport.endpoint.base_url, params.mapped_model?, params.request_query, + gemini_embedding_batch, ), "doubao:embedding" => build_passthrough_path_url( &transport.endpoint.base_url, @@ -287,7 +322,10 @@ fn build_transport_hook_url( None } -fn build_path_params(params: TransportRequestUrlParams<'_>) -> BTreeMap<&'static str, &str> { +fn build_path_params( + params: TransportRequestUrlParams<'_>, + gemini_embedding_batch: bool, +) -> BTreeMap<&'static str, &str> { let mut path_params = BTreeMap::new(); if let Some(model) = params .mapped_model @@ -302,7 +340,11 @@ fn build_path_params(params: TransportRequestUrlParams<'_>) -> BTreeMap<&'static path_params.insert( "action", if provider_api_format == "gemini:embedding" { - "embedContent" + if gemini_embedding_batch { + "batchEmbedContents" + } else { + "embedContent" + } } else if params.upstream_is_stream { "streamGenerateContent" } else { @@ -313,6 +355,14 @@ fn build_path_params(params: TransportRequestUrlParams<'_>) -> BTreeMap<&'static path_params } +fn normalize_gemini_embedding_action_path(path: &str, batch: bool) -> String { + if batch { + path.replace(":embedContent", ":batchEmbedContents") + } else { + path.replace(":batchEmbedContents", ":embedContent") + } +} + fn build_provider_embedding_v1_url(upstream_base_url: &str, query: Option<&str>) -> Option { build_provider_v1_url(upstream_base_url, "/embeddings", "/v1/embeddings", query) } @@ -345,6 +395,7 @@ fn build_gemini_embedding_url( upstream_base_url: &str, model: &str, query: Option<&str>, + batch: bool, ) -> Option { let trimmed_base_url = upstream_base_url .trim() @@ -357,12 +408,17 @@ fn build_gemini_embedding_url( return None; } - let path = if trimmed_base_url.ends_with("/v1beta") { - format!("/models/{trimmed_model}:embedContent") - } else if trimmed_base_url.contains("/v1beta/models/") { - ":embedContent".to_string() + let action = if batch { + "batchEmbedContents" } else { - format!("/v1beta/models/{trimmed_model}:embedContent") + "embedContent" + }; + let path = if trimmed_base_url.ends_with("/v1beta") { + format!("/models/{trimmed_model}:{action}") + } else if trimmed_base_url.contains("/v1beta/models/") { + format!(":{action}") + } else { + format!("/v1beta/models/{trimmed_model}:{action}") }; build_passthrough_path_url(upstream_base_url, &path, query, &["key"]) } @@ -440,12 +496,13 @@ fn custom_path_template_regex() -> &'static Regex { mod tests { use super::{ build_kiro_cross_format_upstream_url, build_transport_request_url, - TransportRequestUrlParams, + build_transport_request_url_for_request_body, TransportRequestUrlParams, }; use crate::snapshot::{ GatewayProviderTransportEndpoint, GatewayProviderTransportKey, GatewayProviderTransportProvider, GatewayProviderTransportSnapshot, }; + use serde_json::json; fn sample_transport( provider_type: &str, @@ -829,6 +886,82 @@ mod tests { ); } + #[test] + fn gemini_embedding_batch_body_uses_batch_endpoint() { + let gemini = sample_transport( + "gemini", + "gemini:embedding", + "https://generativelanguage.googleapis.com/v1beta", + None, + ); + let batch_body = json!({ + "requests": [ + { + "model": "models/gemini-embedding-001", + "content": {"parts": [{"text": "alpha"}]} + }, + { + "model": "models/gemini-embedding-001", + "content": {"parts": [{"text": "beta"}]} + } + ] + }); + + assert_eq!( + build_transport_request_url_for_request_body( + &gemini, + TransportRequestUrlParams { + provider_api_format: "gemini:embedding", + mapped_model: Some("gemini-embedding-001"), + upstream_is_stream: false, + request_query: Some("key=client-key&foo=bar"), + kiro_api_region: None, + }, + Some(&batch_body), + ) + .as_deref(), + Some( + "https://generativelanguage.googleapis.com/v1beta/models/gemini-embedding-001:batchEmbedContents?foo=bar" + ) + ); + } + + #[test] + fn gemini_embedding_custom_action_template_follows_batch_body() { + let gemini = sample_transport( + "gemini", + "gemini:embedding", + "https://generativelanguage.googleapis.com", + Some("/v1beta/models/{model}:{action}"), + ); + let batch_body = json!({ + "requests": [ + { + "model": "models/gemini-embedding-001", + "content": {"parts": [{"text": "alpha"}]} + } + ] + }); + + assert_eq!( + build_transport_request_url_for_request_body( + &gemini, + TransportRequestUrlParams { + provider_api_format: "gemini:embedding", + mapped_model: Some("gemini-embedding-001"), + upstream_is_stream: false, + request_query: None, + kiro_api_region: None, + }, + Some(&batch_body), + ) + .as_deref(), + Some( + "https://generativelanguage.googleapis.com/v1beta/models/gemini-embedding-001:batchEmbedContents" + ) + ); + } + #[test] fn rerank_request_url_builds_provider_default_paths() { let openai = sample_transport( diff --git a/crates/aether-provider-transport/src/same_format_provider/mod.rs b/crates/aether-provider-transport/src/same_format_provider/mod.rs index 20b8a1c13..90a416160 100644 --- a/crates/aether-provider-transport/src/same_format_provider/mod.rs +++ b/crates/aether-provider-transport/src/same_format_provider/mod.rs @@ -26,7 +26,10 @@ use crate::vertex::{ is_vertex_service_account_transport_context, is_vertex_transport_context, local_vertex_gemini_transport_unsupported_reason_with_network, }; -use crate::{build_transport_request_url, ensure_upstream_auth_header, TransportRequestUrlParams}; +use crate::{ + build_transport_request_url_for_request_body, ensure_upstream_auth_header, + TransportRequestUrlParams, +}; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SameFormatProviderFamily { @@ -76,6 +79,7 @@ pub struct SameFormatProviderUpstreamUrlParams<'a> { pub upstream_is_stream: bool, pub request_query: Option<&'a str>, pub kiro_api_region: Option<&'a str>, + pub provider_request_body: Option<&'a Value>, } #[derive(Debug, Clone, Copy)] @@ -256,7 +260,7 @@ pub fn build_same_format_provider_upstream_url( transport: &GatewayProviderTransportSnapshot, params: SameFormatProviderUpstreamUrlParams<'_>, ) -> Option { - build_transport_request_url( + build_transport_request_url_for_request_body( transport, TransportRequestUrlParams { provider_api_format: params.provider_api_format, @@ -265,6 +269,7 @@ pub fn build_same_format_provider_upstream_url( request_query: params.request_query, kiro_api_region: params.kiro_api_region, }, + params.provider_request_body, ) } diff --git a/crates/aether-usage-runtime/src/report.rs b/crates/aether-usage-runtime/src/report.rs index e619ff884..08fe03d4a 100644 --- a/crates/aether-usage-runtime/src/report.rs +++ b/crates/aether-usage-runtime/src/report.rs @@ -234,6 +234,9 @@ pub fn is_local_ai_sync_report_kind(report_kind: &str) -> bool { | "openai_cli_sync_success" | "openai_image_sync_success" | "openai_image_sync_error" + | "openai_embedding_sync_success" + | "openai_embedding_sync_error" + | "gemini_embedding_sync_success" | "claude_cli_sync_success" | "gemini_cli_sync_success" | "openai_cli_sync_error" @@ -426,6 +429,13 @@ mod tests { )); assert!(is_local_ai_sync_report_kind("openai_image_sync_success")); assert!(is_local_ai_sync_report_kind("openai_image_sync_error")); + assert!(is_local_ai_sync_report_kind( + "openai_embedding_sync_success" + )); + assert!(is_local_ai_sync_report_kind("openai_embedding_sync_error")); + assert!(is_local_ai_sync_report_kind( + "gemini_embedding_sync_success" + )); assert!(is_local_ai_sync_report_kind("gemini_files_delete_mapping")); assert!(!is_local_ai_sync_report_kind("unknown_sync_kind")); }