mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-01 17:00:21 +08:00
fix(provider): 将rust分支的gemini cli端点行为对齐到python分支 (#321)
* fix(provider): 对齐 Vertex/Gemini 上游发包与 Python master - provider-transport: 为 custom+aiplatform 推断 Vertex API key 上下文并统一 URL 构建顺序,复用共享 request_url 构建最终上游地址 - ai-pipeline/gateway: Vertex Gemini 路径改为仅使用 URL query key,不再向上游附带 x-goog-api-key header;同步对齐 standard/admin/test-connection/runtime miss 摘要中的最终 URL - gemini conversion: 按 Python master 输出 Gemini 请求体,补齐 system_instruction / generation_config / tool_config / function_declarations 形态,并移植 Gemini schema 清洗逻辑 - scheduler/executor: 将最终 upstream_url、mapped_model、key_name 写入候选 extra_data,运行时 miss 诊断优先展示真实展开后的上游 URL 便于服务器排障 * fix(provider): 修复 Vertex provider 测试与本地调度链路 * fix(provider): 对齐 Vertex 本地执行与 Rust CI
This commit is contained in:
@@ -8,6 +8,11 @@ use aether_provider_transport::policy::{
|
||||
local_openai_chat_transport_unsupported_reason,
|
||||
local_standard_transport_unsupported_reason_with_network,
|
||||
};
|
||||
use aether_provider_transport::vertex::{
|
||||
is_vertex_api_key_transport_context,
|
||||
local_vertex_api_key_gemini_transport_unsupported_reason_with_network,
|
||||
resolve_local_vertex_api_key_query_auth, VERTEX_API_KEY_QUERY_PARAM,
|
||||
};
|
||||
use aether_provider_transport::GatewayProviderTransportSnapshot;
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -255,6 +260,9 @@ pub fn request_conversion_transport_unsupported_reason(
|
||||
"claude:cli" => {
|
||||
local_standard_transport_unsupported_reason_with_network(transport, "claude:cli")
|
||||
}
|
||||
"gemini:chat" | "gemini:cli" if is_vertex_api_key_transport_context(transport) => {
|
||||
local_vertex_api_key_gemini_transport_unsupported_reason_with_network(transport)
|
||||
}
|
||||
"gemini:chat" => {
|
||||
local_gemini_transport_unsupported_reason_with_network(transport, "gemini:chat")
|
||||
}
|
||||
@@ -279,7 +287,14 @@ pub fn request_conversion_direct_auth(
|
||||
"openai:chat" | "openai:cli" | "openai:compact" => {
|
||||
resolve_local_openai_bearer_auth(transport)
|
||||
}
|
||||
"gemini:chat" | "gemini:cli" => resolve_local_gemini_auth(transport),
|
||||
"gemini:chat" | "gemini:cli" => {
|
||||
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))
|
||||
} else {
|
||||
resolve_local_gemini_auth(transport)
|
||||
}
|
||||
}
|
||||
"claude:chat" | "claude:cli" => resolve_local_standard_auth(transport),
|
||||
_ => None,
|
||||
}
|
||||
@@ -729,4 +744,67 @@ mod tests {
|
||||
"openai:cli"
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn vertex_gemini_transport_supports_cross_format_conversion_with_query_auth() {
|
||||
let transport = GatewayProviderTransportSnapshot {
|
||||
provider: GatewayProviderTransportProvider {
|
||||
id: "provider-vertex".to_string(),
|
||||
name: "vertex".to_string(),
|
||||
provider_type: "vertex_ai".to_string(),
|
||||
website: None,
|
||||
is_active: true,
|
||||
keep_priority_on_conversion: false,
|
||||
enable_format_conversion: true,
|
||||
concurrent_limit: None,
|
||||
max_retries: None,
|
||||
proxy: None,
|
||||
request_timeout_secs: None,
|
||||
stream_first_byte_timeout_secs: None,
|
||||
config: None,
|
||||
},
|
||||
endpoint: GatewayProviderTransportEndpoint {
|
||||
id: "endpoint-vertex".to_string(),
|
||||
provider_id: "provider-vertex".to_string(),
|
||||
api_format: "gemini:chat".to_string(),
|
||||
api_family: Some("gemini".to_string()),
|
||||
endpoint_kind: Some("chat".to_string()),
|
||||
is_active: true,
|
||||
base_url: "https://aiplatform.googleapis.com".to_string(),
|
||||
header_rules: None,
|
||||
body_rules: None,
|
||||
max_retries: None,
|
||||
custom_path: None,
|
||||
config: None,
|
||||
format_acceptance_config: None,
|
||||
proxy: None,
|
||||
},
|
||||
key: GatewayProviderTransportKey {
|
||||
id: "key-vertex".to_string(),
|
||||
provider_id: "provider-vertex".to_string(),
|
||||
name: "key".to_string(),
|
||||
auth_type: "api_key".to_string(),
|
||||
is_active: true,
|
||||
api_formats: Some(vec!["gemini:chat".to_string()]),
|
||||
allowed_models: None,
|
||||
capabilities: None,
|
||||
rate_multipliers: None,
|
||||
global_priority_by_format: None,
|
||||
expires_at_unix_secs: None,
|
||||
proxy: None,
|
||||
fingerprint: None,
|
||||
decrypted_api_key: "vertex-secret".to_string(),
|
||||
decrypted_auth_config: None,
|
||||
},
|
||||
};
|
||||
|
||||
assert!(request_conversion_transport_supported(
|
||||
&transport,
|
||||
RequestConversionKind::ToGeminiStandard
|
||||
));
|
||||
assert_eq!(
|
||||
request_conversion_direct_auth(&transport, RequestConversionKind::ToGeminiStandard),
|
||||
Some(("key".to_string(), "vertex-secret".to_string()))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@ use crate::planner::openai::{
|
||||
pub fn convert_openai_chat_request_to_gemini_request(
|
||||
body_json: &Value,
|
||||
mapped_model: &str,
|
||||
upstream_is_stream: bool,
|
||||
_upstream_is_stream: bool,
|
||||
) -> Option<Value> {
|
||||
let request = body_json.as_object()?;
|
||||
let mut system_segments = Vec::new();
|
||||
@@ -123,14 +123,16 @@ pub fn convert_openai_chat_request_to_gemini_request(
|
||||
}
|
||||
|
||||
let mut output = Map::new();
|
||||
output.insert("model".to_string(), Value::String(mapped_model.to_string()));
|
||||
if !mapped_model.trim().is_empty() {
|
||||
output.insert(
|
||||
"model".to_string(),
|
||||
Value::String(mapped_model.trim().to_string()),
|
||||
);
|
||||
}
|
||||
output.insert(
|
||||
"contents".to_string(),
|
||||
Value::Array(compact_gemini_contents(contents)),
|
||||
);
|
||||
if upstream_is_stream {
|
||||
output.insert("stream".to_string(), Value::Bool(true));
|
||||
}
|
||||
let system_text = system_segments
|
||||
.into_iter()
|
||||
.filter(|value| !value.trim().is_empty())
|
||||
@@ -203,6 +205,8 @@ pub fn convert_openai_chat_request_to_gemini_request(
|
||||
.and_then(|json_schema| json_schema.get("schema"))
|
||||
.cloned()
|
||||
{
|
||||
let mut schema = schema;
|
||||
clean_gemini_schema(&mut schema);
|
||||
generation_config.insert("responseSchema".to_string(), schema);
|
||||
}
|
||||
}
|
||||
@@ -232,18 +236,17 @@ pub fn convert_openai_chat_request_to_gemini_request(
|
||||
}
|
||||
if let Some(extra_body) = request.get("extra_body").and_then(Value::as_object) {
|
||||
if let Some(google) = extra_body.get("google").and_then(Value::as_object) {
|
||||
if let Some(existing) = output
|
||||
.get_mut("generationConfig")
|
||||
.and_then(Value::as_object_mut)
|
||||
{
|
||||
if let Some(response_modalities) = google.get("response_modalities").cloned() {
|
||||
existing.insert("responseModalities".to_string(), response_modalities);
|
||||
}
|
||||
if let Some(thinking_config) = google.get("thinking_config").cloned() {
|
||||
existing
|
||||
.entry("thinkingConfig".to_string())
|
||||
.or_insert(thinking_config);
|
||||
}
|
||||
let existing = output
|
||||
.entry("generationConfig".to_string())
|
||||
.or_insert_with(|| Value::Object(Map::new()))
|
||||
.as_object_mut()?;
|
||||
if let Some(response_modalities) = google.get("response_modalities").cloned() {
|
||||
existing.insert("responseModalities".to_string(), response_modalities);
|
||||
}
|
||||
if let Some(thinking_config) = google.get("thinking_config").cloned() {
|
||||
existing
|
||||
.entry("thinkingConfig".to_string())
|
||||
.or_insert(thinking_config);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -394,6 +397,10 @@ fn convert_openai_tools_to_gemini(
|
||||
function
|
||||
.get("parameters")
|
||||
.cloned()
|
||||
.map(|mut schema| {
|
||||
clean_gemini_schema(&mut schema);
|
||||
schema
|
||||
})
|
||||
.unwrap_or_else(|| json!({})),
|
||||
);
|
||||
declarations.push(Value::Object(declaration));
|
||||
@@ -483,6 +490,517 @@ fn compact_gemini_contents(contents: Vec<Value>) -> Vec<Value> {
|
||||
compact
|
||||
}
|
||||
|
||||
const ALLOWED_SCHEMA_FIELDS: &[&str] = &[
|
||||
"type",
|
||||
"description",
|
||||
"properties",
|
||||
"required",
|
||||
"items",
|
||||
"enum",
|
||||
"title",
|
||||
];
|
||||
|
||||
const CONSTRAINT_FIELDS: &[(&str, &str)] = &[
|
||||
("minLength", "minLen"),
|
||||
("maxLength", "maxLen"),
|
||||
("pattern", "pattern"),
|
||||
("minimum", "min"),
|
||||
("maximum", "max"),
|
||||
("multipleOf", "multipleOf"),
|
||||
("exclusiveMinimum", "exclMin"),
|
||||
("exclusiveMaximum", "exclMax"),
|
||||
("minItems", "minItems"),
|
||||
("maxItems", "maxItems"),
|
||||
("format", "format"),
|
||||
];
|
||||
|
||||
fn clean_gemini_schema(value: &mut Value) {
|
||||
if !value.is_object() {
|
||||
return;
|
||||
}
|
||||
|
||||
let mut defs = Map::new();
|
||||
collect_all_defs(value, &mut defs);
|
||||
if let Some(object) = value.as_object_mut() {
|
||||
object.remove("$defs");
|
||||
object.remove("definitions");
|
||||
}
|
||||
let mut seen = Vec::new();
|
||||
flatten_refs(value, &defs, &mut seen);
|
||||
clean_schema_recursive(value, true);
|
||||
}
|
||||
|
||||
fn collect_all_defs(value: &Value, defs: &mut Map<String, Value>) {
|
||||
match value {
|
||||
Value::Object(object) => {
|
||||
for defs_key in ["$defs", "definitions"] {
|
||||
if let Some(Value::Object(inner_defs)) = object.get(defs_key) {
|
||||
for (key, inner) in inner_defs {
|
||||
defs.entry(key.clone()).or_insert_with(|| inner.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
for (key, inner) in object {
|
||||
if key != "$defs" && key != "definitions" {
|
||||
collect_all_defs(inner, defs);
|
||||
}
|
||||
}
|
||||
}
|
||||
Value::Array(items) => {
|
||||
for item in items {
|
||||
collect_all_defs(item, defs);
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
fn flatten_refs(value: &mut Value, defs: &Map<String, Value>, seen: &mut Vec<String>) {
|
||||
match value {
|
||||
Value::Object(object) => {
|
||||
let ref_path = object
|
||||
.remove("$ref")
|
||||
.and_then(|value| value.as_str().map(ToOwned::to_owned));
|
||||
if let Some(ref_path) = ref_path {
|
||||
let ref_name = ref_path.rsplit('/').next().unwrap_or_default().to_string();
|
||||
if seen.iter().any(|value| value == &ref_name) {
|
||||
object
|
||||
.entry("type".to_string())
|
||||
.or_insert_with(|| Value::String("string".to_string()));
|
||||
append_schema_hint(object, &format!("(Circular $ref: {ref_path})"));
|
||||
return;
|
||||
}
|
||||
seen.push(ref_name.clone());
|
||||
if let Some(Value::Object(def_schema)) = defs.get(&ref_name) {
|
||||
for (key, inner) in def_schema {
|
||||
if !object.contains_key(key) {
|
||||
object.insert(key.clone(), inner.clone());
|
||||
}
|
||||
}
|
||||
flatten_refs(value, defs, seen);
|
||||
} else {
|
||||
object
|
||||
.entry("type".to_string())
|
||||
.or_insert_with(|| Value::String("string".to_string()));
|
||||
append_schema_hint(object, &format!("(Unresolved $ref: {ref_path})"));
|
||||
}
|
||||
seen.pop();
|
||||
return;
|
||||
}
|
||||
for inner in object.values_mut() {
|
||||
flatten_refs(inner, defs, seen);
|
||||
}
|
||||
}
|
||||
Value::Array(items) => {
|
||||
for item in items {
|
||||
flatten_refs(item, defs, seen);
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
fn clean_schema_recursive(value: &mut Value, is_schema_node: bool) -> bool {
|
||||
let Some(object) = value.as_object_mut() else {
|
||||
if let Some(items) = value.as_array_mut() {
|
||||
for item in items {
|
||||
clean_schema_recursive(item, is_schema_node);
|
||||
}
|
||||
}
|
||||
return false;
|
||||
};
|
||||
|
||||
let mut is_nullable = false;
|
||||
merge_all_of(object);
|
||||
|
||||
if (object.get("type").and_then(Value::as_str) == Some("object")
|
||||
|| object.contains_key("properties"))
|
||||
&& object.contains_key("items")
|
||||
{
|
||||
let items = object.remove("items");
|
||||
if let Some(Value::Object(items)) = items {
|
||||
let props = object
|
||||
.entry("properties".to_string())
|
||||
.or_insert_with(|| Value::Object(Map::new()))
|
||||
.as_object_mut();
|
||||
if let Some(props) = props {
|
||||
for (key, inner) in items {
|
||||
props.entry(key).or_insert(inner);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut nullable_keys = Vec::new();
|
||||
if let Some(Value::Object(props)) = object.get_mut("properties") {
|
||||
for (key, inner) in props.iter_mut() {
|
||||
if clean_schema_recursive(inner, true) {
|
||||
nullable_keys.push(key.clone());
|
||||
}
|
||||
}
|
||||
if !object.contains_key("type") {
|
||||
object.insert("type".to_string(), Value::String("object".to_string()));
|
||||
}
|
||||
}
|
||||
if !nullable_keys.is_empty() {
|
||||
if let Some(Value::Array(required)) = object.get_mut("required") {
|
||||
required.retain(|item| {
|
||||
item.as_str()
|
||||
.is_some_and(|value| !nullable_keys.iter().any(|candidate| candidate == value))
|
||||
});
|
||||
if required.is_empty() {
|
||||
object.remove("required");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(items) = object.get_mut("items") {
|
||||
if items.is_object() {
|
||||
clean_schema_recursive(items, true);
|
||||
if !object.contains_key("type") {
|
||||
object.insert("type".to_string(), Value::String("array".to_string()));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if !object.contains_key("properties") && !object.contains_key("items") {
|
||||
for (key, inner) in object.iter_mut() {
|
||||
if !matches!(key.as_str(), "anyOf" | "oneOf" | "allOf" | "enum" | "type")
|
||||
&& (inner.is_object() || inner.is_array())
|
||||
{
|
||||
clean_schema_recursive(inner, false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for combo_key in ["anyOf", "oneOf"] {
|
||||
if let Some(Value::Array(combo)) = object.get_mut(combo_key) {
|
||||
for branch in combo {
|
||||
if branch.is_object() {
|
||||
clean_schema_recursive(branch, true);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let should_merge_union = object.get("type").is_none()
|
||||
|| object.get("type").and_then(Value::as_str) == Some("object");
|
||||
if should_merge_union {
|
||||
let union = object
|
||||
.get("anyOf")
|
||||
.and_then(Value::as_array)
|
||||
.or_else(|| object.get("oneOf").and_then(Value::as_array))
|
||||
.cloned();
|
||||
if let Some(union) = union {
|
||||
let (best, all_types) = extract_best_schema_branch(&union);
|
||||
if let Some(Value::Object(best_object)) = best {
|
||||
for (key, inner) in best_object {
|
||||
if key == "properties" {
|
||||
let target = object
|
||||
.entry("properties".to_string())
|
||||
.or_insert_with(|| Value::Object(Map::new()))
|
||||
.as_object_mut();
|
||||
if let (Some(target), Value::Object(props)) = (target, inner) {
|
||||
for (prop_key, prop_value) in props {
|
||||
target.entry(prop_key).or_insert(prop_value);
|
||||
}
|
||||
}
|
||||
} else if key == "required" {
|
||||
let target = object
|
||||
.entry("required".to_string())
|
||||
.or_insert_with(|| Value::Array(Vec::new()))
|
||||
.as_array_mut();
|
||||
if let (Some(target), Value::Array(required)) = (target, inner) {
|
||||
for required_value in required {
|
||||
if !target.iter().any(|value| value == &required_value) {
|
||||
target.push(required_value);
|
||||
}
|
||||
}
|
||||
}
|
||||
} else if !object.contains_key(&key) {
|
||||
object.insert(key, inner);
|
||||
}
|
||||
}
|
||||
if all_types.len() > 1 {
|
||||
append_schema_hint(object, &format!("Accepts: {}", all_types.join(" | ")));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
object.remove("anyOf");
|
||||
object.remove("oneOf");
|
||||
|
||||
let is_not_schema_payload = object.contains_key("functionCall")
|
||||
|| object.contains_key("functionResponse")
|
||||
|| object.contains_key("function_call")
|
||||
|| object.contains_key("function_response");
|
||||
let has_standard = object
|
||||
.keys()
|
||||
.any(|key| ALLOWED_SCHEMA_FIELDS.iter().any(|allowed| key == allowed));
|
||||
|
||||
if is_schema_node && !has_standard && !object.is_empty() && !is_not_schema_payload {
|
||||
let keys = object.keys().cloned().collect::<Vec<_>>();
|
||||
let mut new_props = Map::new();
|
||||
for key in keys {
|
||||
if let Some(inner) = object.remove(&key) {
|
||||
new_props.insert(key, inner);
|
||||
}
|
||||
}
|
||||
for inner in new_props.values_mut() {
|
||||
if inner.is_object() {
|
||||
clean_schema_recursive(inner, true);
|
||||
}
|
||||
}
|
||||
object.insert("type".to_string(), Value::String("object".to_string()));
|
||||
object.insert("properties".to_string(), Value::Object(new_props));
|
||||
}
|
||||
|
||||
let looks_like_schema = (is_schema_node || has_standard || object.contains_key("properties"))
|
||||
&& !is_not_schema_payload;
|
||||
if looks_like_schema {
|
||||
move_constraints_to_description(object);
|
||||
let keys_to_remove = object
|
||||
.keys()
|
||||
.filter(|key| !ALLOWED_SCHEMA_FIELDS.iter().any(|allowed| *key == allowed))
|
||||
.cloned()
|
||||
.collect::<Vec<_>>();
|
||||
for key in keys_to_remove {
|
||||
object.remove(&key);
|
||||
}
|
||||
|
||||
if object.get("type").and_then(Value::as_str) == Some("object")
|
||||
&& !object.contains_key("properties")
|
||||
{
|
||||
object.insert("properties".to_string(), Value::Object(Map::new()));
|
||||
}
|
||||
|
||||
let valid_keys = object
|
||||
.get("properties")
|
||||
.and_then(Value::as_object)
|
||||
.map(|props| props.keys().cloned().collect::<Vec<_>>())
|
||||
.unwrap_or_default();
|
||||
if let Some(Value::Array(required)) = object.get_mut("required") {
|
||||
required.retain(|item| {
|
||||
item.as_str()
|
||||
.is_some_and(|value| valid_keys.iter().any(|candidate| candidate == value))
|
||||
});
|
||||
if required.is_empty() {
|
||||
object.remove("required");
|
||||
}
|
||||
}
|
||||
|
||||
if !object.contains_key("type") {
|
||||
let inferred_type = if object.contains_key("enum") {
|
||||
"string"
|
||||
} else if object.contains_key("properties") {
|
||||
"object"
|
||||
} else if object.contains_key("items") {
|
||||
"array"
|
||||
} else {
|
||||
"string"
|
||||
};
|
||||
object.insert("type".to_string(), Value::String(inferred_type.to_string()));
|
||||
}
|
||||
|
||||
let fallback_type = if object.contains_key("properties") {
|
||||
"object"
|
||||
} else if object.contains_key("items") {
|
||||
"array"
|
||||
} else {
|
||||
"string"
|
||||
};
|
||||
let selected_type = match object.get("type") {
|
||||
Some(Value::String(type_name)) => {
|
||||
let lower = type_name.to_ascii_lowercase();
|
||||
if lower == "null" {
|
||||
is_nullable = true;
|
||||
None
|
||||
} else {
|
||||
Some(lower)
|
||||
}
|
||||
}
|
||||
Some(Value::Array(types)) => {
|
||||
let mut selected = None;
|
||||
for item in types {
|
||||
if let Some(type_name) = item.as_str() {
|
||||
let lower = type_name.to_ascii_lowercase();
|
||||
if lower == "null" {
|
||||
is_nullable = true;
|
||||
} else if selected.is_none() {
|
||||
selected = Some(lower);
|
||||
}
|
||||
}
|
||||
}
|
||||
selected
|
||||
}
|
||||
_ => None,
|
||||
};
|
||||
object.insert(
|
||||
"type".to_string(),
|
||||
Value::String(selected_type.unwrap_or_else(|| fallback_type.to_string())),
|
||||
);
|
||||
|
||||
if is_nullable {
|
||||
append_schema_hint(object, "(nullable)");
|
||||
}
|
||||
|
||||
if let Some(Value::Array(items)) = object.get_mut("enum") {
|
||||
for item in items.iter_mut() {
|
||||
if !item.is_string() {
|
||||
*item = Value::String(match item {
|
||||
Value::Null => "null".to_string(),
|
||||
_ => item.to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
is_nullable
|
||||
}
|
||||
|
||||
fn merge_all_of(object: &mut Map<String, Value>) {
|
||||
let all_of = object.remove("allOf");
|
||||
let Some(Value::Array(all_of)) = all_of else {
|
||||
return;
|
||||
};
|
||||
|
||||
let mut merged_props = Map::new();
|
||||
let mut merged_required = Vec::new();
|
||||
let mut other_fields = Map::new();
|
||||
|
||||
for item in all_of {
|
||||
let Value::Object(item) = item else {
|
||||
continue;
|
||||
};
|
||||
if let Some(Value::Object(props)) = item.get("properties") {
|
||||
for (key, value) in props {
|
||||
merged_props.insert(key.clone(), value.clone());
|
||||
}
|
||||
}
|
||||
if let Some(Value::Array(required)) = item.get("required") {
|
||||
for value in required {
|
||||
if !merged_required.iter().any(|existing| existing == value) {
|
||||
merged_required.push(value.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
for (key, value) in item {
|
||||
if !matches!(key.as_str(), "properties" | "required" | "allOf")
|
||||
&& !other_fields.contains_key(&key)
|
||||
{
|
||||
other_fields.insert(key, value);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for (key, value) in other_fields {
|
||||
object.entry(key).or_insert(value);
|
||||
}
|
||||
if !merged_props.is_empty() {
|
||||
let target = object
|
||||
.entry("properties".to_string())
|
||||
.or_insert_with(|| Value::Object(Map::new()))
|
||||
.as_object_mut();
|
||||
if let Some(target) = target {
|
||||
for (key, value) in merged_props {
|
||||
target.entry(key).or_insert(value);
|
||||
}
|
||||
}
|
||||
}
|
||||
if !merged_required.is_empty() {
|
||||
let target = object
|
||||
.entry("required".to_string())
|
||||
.or_insert_with(|| Value::Array(Vec::new()))
|
||||
.as_array_mut();
|
||||
if let Some(target) = target {
|
||||
for value in merged_required {
|
||||
if !target.iter().any(|existing| existing == &value) {
|
||||
target.push(value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn extract_best_schema_branch(union: &[Value]) -> (Option<Value>, Vec<String>) {
|
||||
let mut best = None;
|
||||
let mut best_score = -1;
|
||||
let mut all_types = Vec::new();
|
||||
|
||||
for item in union {
|
||||
let score = score_schema_branch(item);
|
||||
if let Some(type_name) = schema_type_name(item) {
|
||||
if !all_types.iter().any(|existing| existing == type_name) {
|
||||
all_types.push(type_name.to_string());
|
||||
}
|
||||
}
|
||||
if score > best_score {
|
||||
best_score = score;
|
||||
best = Some(item.clone());
|
||||
}
|
||||
}
|
||||
|
||||
(best, all_types)
|
||||
}
|
||||
|
||||
fn score_schema_branch(value: &Value) -> i32 {
|
||||
let Some(object) = value.as_object() else {
|
||||
return 0;
|
||||
};
|
||||
if object.contains_key("properties")
|
||||
|| object.get("type").and_then(Value::as_str) == Some("object")
|
||||
{
|
||||
return 3;
|
||||
}
|
||||
if object.contains_key("items") || object.get("type").and_then(Value::as_str) == Some("array") {
|
||||
return 2;
|
||||
}
|
||||
if object
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
.is_some_and(|value| value != "null")
|
||||
{
|
||||
return 1;
|
||||
}
|
||||
0
|
||||
}
|
||||
|
||||
fn schema_type_name(value: &Value) -> Option<&str> {
|
||||
let object = value.as_object()?;
|
||||
object
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
.or_else(|| object.contains_key("properties").then_some("object"))
|
||||
.or_else(|| object.contains_key("items").then_some("array"))
|
||||
}
|
||||
|
||||
fn move_constraints_to_description(object: &mut Map<String, Value>) {
|
||||
let hints = CONSTRAINT_FIELDS
|
||||
.iter()
|
||||
.filter_map(|(field, label)| object.get(*field).map(|value| format!("{label}: {value}")))
|
||||
.collect::<Vec<_>>();
|
||||
if !hints.is_empty() {
|
||||
append_schema_hint(object, &format!("[Constraint: {}]", hints.join(", ")));
|
||||
}
|
||||
}
|
||||
|
||||
fn append_schema_hint(object: &mut Map<String, Value>, hint: &str) {
|
||||
let existing = object
|
||||
.get("description")
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or_default();
|
||||
if existing.contains(hint) {
|
||||
return;
|
||||
}
|
||||
let next = if existing.trim().is_empty() {
|
||||
hint.to_string()
|
||||
} else {
|
||||
format!("{existing} {hint}")
|
||||
};
|
||||
object.insert("description".to_string(), Value::String(next));
|
||||
}
|
||||
|
||||
fn parse_data_url(value: &str) -> Option<(String, String)> {
|
||||
let rest = value.strip_prefix("data:")?;
|
||||
let (meta, data) = rest.split_once(",")?;
|
||||
@@ -572,6 +1090,7 @@ mod tests {
|
||||
convert_openai_chat_request_to_gemini_request(&request, "gemini-2.5-pro", false)
|
||||
.expect("request should convert");
|
||||
|
||||
assert_eq!(converted["model"], "gemini-2.5-pro");
|
||||
assert_eq!(converted["generationConfig"]["seed"], 7);
|
||||
assert_eq!(converted["tools"][0], json!({ "codeExecution": {} }));
|
||||
assert_eq!(converted["tools"][1], json!({ "googleSearch": {} }));
|
||||
|
||||
@@ -1,10 +1,9 @@
|
||||
use std::borrow::Cow;
|
||||
|
||||
use aether_provider_transport::url::{
|
||||
build_claude_messages_url, build_gemini_content_url, build_openai_chat_url,
|
||||
build_openai_cli_url, build_passthrough_path_url,
|
||||
use aether_provider_transport::{
|
||||
apply_local_body_rules, build_transport_request_url, GatewayProviderTransportSnapshot,
|
||||
TransportRequestUrlParams,
|
||||
};
|
||||
use aether_provider_transport::{apply_local_body_rules, GatewayProviderTransportSnapshot};
|
||||
use serde_json::Value;
|
||||
|
||||
use crate::conversion::request::{
|
||||
@@ -136,45 +135,16 @@ pub fn build_standard_upstream_url(
|
||||
provider_api_format: &str,
|
||||
upstream_is_stream: bool,
|
||||
) -> Option<String> {
|
||||
let custom_path = transport
|
||||
.endpoint
|
||||
.custom_path
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty());
|
||||
|
||||
match custom_path {
|
||||
Some(path) => {
|
||||
build_passthrough_path_url(&transport.endpoint.base_url, path, parts.uri.query(), &[])
|
||||
}
|
||||
None => match provider_api_format.trim().to_ascii_lowercase().as_str() {
|
||||
"openai:chat" => Some(build_openai_chat_url(
|
||||
&transport.endpoint.base_url,
|
||||
parts.uri.query(),
|
||||
)),
|
||||
"openai:cli" => Some(build_openai_cli_url(
|
||||
&transport.endpoint.base_url,
|
||||
parts.uri.query(),
|
||||
false,
|
||||
)),
|
||||
"openai:compact" => Some(build_openai_cli_url(
|
||||
&transport.endpoint.base_url,
|
||||
parts.uri.query(),
|
||||
true,
|
||||
)),
|
||||
"claude:chat" | "claude:cli" => Some(build_claude_messages_url(
|
||||
&transport.endpoint.base_url,
|
||||
parts.uri.query(),
|
||||
)),
|
||||
"gemini:chat" | "gemini:cli" => build_gemini_content_url(
|
||||
&transport.endpoint.base_url,
|
||||
mapped_model,
|
||||
upstream_is_stream,
|
||||
parts.uri.query(),
|
||||
),
|
||||
_ => None,
|
||||
build_transport_request_url(
|
||||
transport,
|
||||
TransportRequestUrlParams {
|
||||
provider_api_format,
|
||||
mapped_model: Some(mapped_model),
|
||||
upstream_is_stream,
|
||||
request_query: parts.uri.query(),
|
||||
kiro_api_region: None,
|
||||
},
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -3,8 +3,8 @@ use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use aether_contracts::{ExecutionPlan, ExecutionResult, RequestBody};
|
||||
use aether_provider_transport::{
|
||||
resolve_transport_execution_timeouts, resolve_transport_tls_profile,
|
||||
GatewayProviderTransportSnapshot,
|
||||
is_vertex_api_key_transport_context, resolve_transport_execution_timeouts,
|
||||
resolve_transport_tls_profile, GatewayProviderTransportSnapshot,
|
||||
};
|
||||
use base64::engine::general_purpose::{STANDARD, URL_SAFE_NO_PAD};
|
||||
use base64::Engine as _;
|
||||
@@ -61,6 +61,10 @@ pub async fn fetch_models_from_transports(
|
||||
return Ok(build_success_outcome(models, None, true));
|
||||
}
|
||||
|
||||
if transports.iter().any(is_vertex_api_key_transport_context) {
|
||||
return fetch_vertex_models(runtime, transports).await;
|
||||
}
|
||||
|
||||
match provider_type.as_str() {
|
||||
"antigravity" => fetch_antigravity_models(runtime, first_transport).await,
|
||||
"vertex_ai" => fetch_vertex_models(runtime, transports).await,
|
||||
@@ -1061,3 +1065,146 @@ impl OutcomeExt for ModelsFetchOutcome {
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use aether_contracts::{ExecutionResult, ResponseBody};
|
||||
use aether_provider_transport::snapshot::{
|
||||
GatewayProviderTransportEndpoint, GatewayProviderTransportKey,
|
||||
GatewayProviderTransportProvider, GatewayProviderTransportSnapshot,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use serde_json::json;
|
||||
|
||||
use crate::fetch_models_from_transports;
|
||||
use crate::transport::ModelFetchTransportRuntime;
|
||||
|
||||
struct TestRuntime {
|
||||
executed_urls: Arc<Mutex<Vec<String>>>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ModelFetchTransportRuntime for TestRuntime {
|
||||
async fn resolve_local_oauth_request_auth(
|
||||
&self,
|
||||
_transport: &GatewayProviderTransportSnapshot,
|
||||
) -> Result<Option<aether_provider_transport::LocalResolvedOAuthRequestAuth>, String>
|
||||
{
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
async fn resolve_model_fetch_proxy(
|
||||
&self,
|
||||
_transport: &GatewayProviderTransportSnapshot,
|
||||
) -> Option<aether_contracts::ProxySnapshot> {
|
||||
None
|
||||
}
|
||||
|
||||
async fn execute_model_fetch_execution_plan(
|
||||
&self,
|
||||
plan: &aether_contracts::ExecutionPlan,
|
||||
) -> Result<ExecutionResult, String> {
|
||||
self.executed_urls
|
||||
.lock()
|
||||
.expect("executed_urls lock")
|
||||
.push(plan.url.clone());
|
||||
Ok(ExecutionResult {
|
||||
request_id: plan.request_id.clone(),
|
||||
candidate_id: plan.candidate_id.clone(),
|
||||
status_code: 200,
|
||||
headers: BTreeMap::new(),
|
||||
body: Some(ResponseBody {
|
||||
json_body: Some(json!({
|
||||
"models": [{
|
||||
"name": "publishers/google/models/gemini-3.1-pro-preview"
|
||||
}]
|
||||
})),
|
||||
body_bytes_b64: None,
|
||||
}),
|
||||
telemetry: None,
|
||||
error: None,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
fn sample_custom_aiplatform_transport() -> GatewayProviderTransportSnapshot {
|
||||
GatewayProviderTransportSnapshot {
|
||||
provider: GatewayProviderTransportProvider {
|
||||
id: "provider-1".to_string(),
|
||||
name: "Vertex".to_string(),
|
||||
provider_type: "custom".to_string(),
|
||||
website: None,
|
||||
is_active: true,
|
||||
keep_priority_on_conversion: false,
|
||||
enable_format_conversion: true,
|
||||
concurrent_limit: None,
|
||||
max_retries: None,
|
||||
proxy: None,
|
||||
request_timeout_secs: None,
|
||||
stream_first_byte_timeout_secs: None,
|
||||
config: None,
|
||||
},
|
||||
endpoint: GatewayProviderTransportEndpoint {
|
||||
id: "endpoint-1".to_string(),
|
||||
provider_id: "provider-1".to_string(),
|
||||
api_format: "gemini:cli".to_string(),
|
||||
api_family: Some("gemini".to_string()),
|
||||
endpoint_kind: Some("cli".to_string()),
|
||||
is_active: true,
|
||||
base_url: "https://aiplatform.googleapis.com".to_string(),
|
||||
header_rules: None,
|
||||
body_rules: None,
|
||||
max_retries: None,
|
||||
custom_path: Some("/v1/publishers/google/models/{model}:{action}".to_string()),
|
||||
config: None,
|
||||
format_acceptance_config: None,
|
||||
proxy: None,
|
||||
},
|
||||
key: GatewayProviderTransportKey {
|
||||
id: "key-1".to_string(),
|
||||
provider_id: "provider-1".to_string(),
|
||||
name: "key".to_string(),
|
||||
auth_type: "api_key".to_string(),
|
||||
is_active: true,
|
||||
api_formats: Some(vec!["gemini:cli".to_string()]),
|
||||
allowed_models: None,
|
||||
capabilities: None,
|
||||
rate_multipliers: None,
|
||||
global_priority_by_format: None,
|
||||
expires_at_unix_secs: None,
|
||||
proxy: None,
|
||||
fingerprint: None,
|
||||
decrypted_api_key: "vertex-secret".to_string(),
|
||||
decrypted_auth_config: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn custom_aiplatform_transport_uses_vertex_models_fetch_path_and_normalizes_chat_format()
|
||||
{
|
||||
let executed_urls = Arc::new(Mutex::new(Vec::new()));
|
||||
let runtime = TestRuntime {
|
||||
executed_urls: Arc::clone(&executed_urls),
|
||||
};
|
||||
let outcome =
|
||||
fetch_models_from_transports(&runtime, &[sample_custom_aiplatform_transport()])
|
||||
.await
|
||||
.expect("models fetch should succeed");
|
||||
|
||||
let urls = executed_urls.lock().expect("executed_urls lock");
|
||||
assert_eq!(
|
||||
urls.as_slice(),
|
||||
&["https://aiplatform.googleapis.com/v1/publishers/google/models?key=vertex-secret&pageSize=100"]
|
||||
);
|
||||
assert_eq!(outcome.fetched_model_ids, vec!["gemini-3.1-pro-preview"]);
|
||||
assert_eq!(outcome.cached_models.len(), 1);
|
||||
assert_eq!(
|
||||
outcome.cached_models[0]["api_formats"][0].as_str(),
|
||||
Some("gemini:chat")
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ mod network;
|
||||
pub mod oauth_refresh;
|
||||
pub mod policy;
|
||||
pub mod provider_types;
|
||||
mod request_url;
|
||||
pub mod rules;
|
||||
pub mod snapshot;
|
||||
pub mod url;
|
||||
@@ -40,6 +41,7 @@ pub use policy::{
|
||||
local_standard_transport_unsupported_reason_with_network, supports_local_gemini_transport,
|
||||
supports_local_gemini_transport_with_network, supports_local_standard_transport,
|
||||
};
|
||||
pub use request_url::{build_transport_request_url, TransportRequestUrlParams};
|
||||
pub use rules::{
|
||||
apply_local_body_rules, apply_local_header_rules, body_rules_are_locally_supported,
|
||||
body_rules_handle_path, header_rules_are_locally_supported,
|
||||
@@ -48,6 +50,7 @@ pub use snapshot::{
|
||||
read_provider_transport_snapshot, GatewayProviderTransportSnapshot,
|
||||
ProviderTransportSnapshotSource,
|
||||
};
|
||||
pub use vertex::{is_vertex_api_key_transport_context, uses_vertex_api_key_query_auth};
|
||||
pub use video::{
|
||||
reconstruct_local_video_task_snapshot, resolve_local_video_task_transport,
|
||||
VideoTaskTransportSnapshotLookup,
|
||||
|
||||
399
crates/aether-provider-transport/src/request_url.rs
Normal file
399
crates/aether-provider-transport/src/request_url.rs
Normal file
@@ -0,0 +1,399 @@
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::OnceLock;
|
||||
|
||||
use regex::Regex;
|
||||
use url::form_urlencoded;
|
||||
|
||||
use crate::antigravity::{build_antigravity_v1internal_url, AntigravityRequestUrlAction};
|
||||
use crate::claude_code::build_claude_code_messages_url;
|
||||
use crate::snapshot::GatewayProviderTransportSnapshot;
|
||||
use crate::url::{
|
||||
build_claude_messages_url, build_gemini_content_url, build_openai_chat_url,
|
||||
build_openai_cli_url, build_passthrough_path_url,
|
||||
};
|
||||
use crate::vertex::{
|
||||
build_vertex_api_key_gemini_content_url, resolve_local_vertex_api_key_query_auth,
|
||||
};
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub struct TransportRequestUrlParams<'a> {
|
||||
pub provider_api_format: &'a str,
|
||||
pub mapped_model: Option<&'a str>,
|
||||
pub upstream_is_stream: bool,
|
||||
pub request_query: Option<&'a str>,
|
||||
pub kiro_api_region: Option<&'a str>,
|
||||
}
|
||||
|
||||
pub fn build_transport_request_url(
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
params: TransportRequestUrlParams<'_>,
|
||||
) -> Option<String> {
|
||||
if let Some(url) = build_transport_hook_url(transport, params) {
|
||||
return Some(url);
|
||||
}
|
||||
|
||||
let provider_api_format = params.provider_api_format.trim().to_ascii_lowercase();
|
||||
let custom_path = transport
|
||||
.endpoint
|
||||
.custom_path
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.map(|path| expand_custom_path_template(path, build_path_params(params)));
|
||||
|
||||
if let Some(path) = custom_path.as_deref() {
|
||||
let blocked_keys = if provider_api_format.starts_with("gemini:") {
|
||||
&["key"][..]
|
||||
} else {
|
||||
&[][..]
|
||||
};
|
||||
let url = build_passthrough_path_url(
|
||||
&transport.endpoint.base_url,
|
||||
path,
|
||||
params.request_query,
|
||||
blocked_keys,
|
||||
)?;
|
||||
return Some(maybe_add_gemini_stream_alt_sse(
|
||||
url,
|
||||
&provider_api_format,
|
||||
params.upstream_is_stream,
|
||||
));
|
||||
}
|
||||
|
||||
let url = match provider_api_format.as_str() {
|
||||
"openai:chat" => Some(build_openai_chat_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.request_query,
|
||||
)),
|
||||
"openai:cli" => Some(build_openai_cli_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.request_query,
|
||||
false,
|
||||
)),
|
||||
"openai:compact" => Some(build_openai_cli_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.request_query,
|
||||
true,
|
||||
)),
|
||||
"claude:chat" | "claude:cli" => Some(build_claude_messages_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.request_query,
|
||||
)),
|
||||
"gemini:chat" | "gemini:cli" => build_gemini_content_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.mapped_model?,
|
||||
params.upstream_is_stream,
|
||||
params.request_query,
|
||||
),
|
||||
_ => None,
|
||||
}?;
|
||||
|
||||
Some(maybe_add_gemini_stream_alt_sse(
|
||||
url,
|
||||
&provider_api_format,
|
||||
params.upstream_is_stream,
|
||||
))
|
||||
}
|
||||
|
||||
fn build_transport_hook_url(
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
params: TransportRequestUrlParams<'_>,
|
||||
) -> Option<String> {
|
||||
if let Some(api_region) = params.kiro_api_region {
|
||||
return crate::kiro::build_kiro_generate_assistant_response_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.request_query,
|
||||
Some(api_region),
|
||||
);
|
||||
}
|
||||
|
||||
if transport
|
||||
.provider
|
||||
.provider_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case("claude_code")
|
||||
{
|
||||
return Some(build_claude_code_messages_url(
|
||||
&transport.endpoint.base_url,
|
||||
params.request_query,
|
||||
));
|
||||
}
|
||||
|
||||
if params
|
||||
.provider_api_format
|
||||
.trim()
|
||||
.to_ascii_lowercase()
|
||||
.starts_with("gemini:")
|
||||
{
|
||||
if let Some(auth) = resolve_local_vertex_api_key_query_auth(transport) {
|
||||
return build_vertex_api_key_gemini_content_url(
|
||||
params.mapped_model?,
|
||||
params.upstream_is_stream,
|
||||
&auth.value,
|
||||
params.request_query,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if transport
|
||||
.provider
|
||||
.provider_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case("antigravity")
|
||||
{
|
||||
let query = params.request_query.map(|raw| {
|
||||
form_urlencoded::parse(raw.as_bytes())
|
||||
.into_owned()
|
||||
.collect::<BTreeMap<String, String>>()
|
||||
});
|
||||
return build_antigravity_v1internal_url(
|
||||
&transport.endpoint.base_url,
|
||||
if params.upstream_is_stream {
|
||||
AntigravityRequestUrlAction::StreamGenerateContent
|
||||
} else {
|
||||
AntigravityRequestUrlAction::GenerateContent
|
||||
},
|
||||
query.as_ref(),
|
||||
);
|
||||
}
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
fn build_path_params(params: TransportRequestUrlParams<'_>) -> BTreeMap<&'static str, &str> {
|
||||
let mut path_params = BTreeMap::new();
|
||||
if let Some(model) = params
|
||||
.mapped_model
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
path_params.insert("model", model);
|
||||
}
|
||||
if params
|
||||
.provider_api_format
|
||||
.trim()
|
||||
.to_ascii_lowercase()
|
||||
.starts_with("gemini:")
|
||||
{
|
||||
path_params.insert(
|
||||
"action",
|
||||
if params.upstream_is_stream {
|
||||
"streamGenerateContent"
|
||||
} else {
|
||||
"generateContent"
|
||||
},
|
||||
);
|
||||
}
|
||||
path_params
|
||||
}
|
||||
|
||||
fn expand_custom_path_template(path: &str, params: BTreeMap<&'static str, &str>) -> String {
|
||||
if params.is_empty() {
|
||||
return path.to_string();
|
||||
}
|
||||
|
||||
let regex = custom_path_template_regex();
|
||||
let mut missing_key = false;
|
||||
let replaced = regex.replace_all(path, |captures: ®ex::Captures<'_>| {
|
||||
let key = captures
|
||||
.get(1)
|
||||
.map(|value| value.as_str())
|
||||
.unwrap_or_default();
|
||||
match params.get(key).copied() {
|
||||
Some(value) => value.to_string(),
|
||||
None => {
|
||||
missing_key = true;
|
||||
captures
|
||||
.get(0)
|
||||
.map(|value| value.as_str().to_string())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
if missing_key {
|
||||
path.to_string()
|
||||
} else {
|
||||
replaced.into_owned()
|
||||
}
|
||||
}
|
||||
|
||||
fn maybe_add_gemini_stream_alt_sse(
|
||||
upstream_url: String,
|
||||
provider_api_format: &str,
|
||||
upstream_is_stream: bool,
|
||||
) -> String {
|
||||
if !provider_api_format.starts_with("gemini:") || !upstream_is_stream {
|
||||
return upstream_url;
|
||||
}
|
||||
|
||||
let has_alt = upstream_url
|
||||
.split_once('?')
|
||||
.map(|(_, query)| {
|
||||
form_urlencoded::parse(query.as_bytes())
|
||||
.any(|(key, _)| key.as_ref().eq_ignore_ascii_case("alt"))
|
||||
})
|
||||
.unwrap_or(false);
|
||||
if has_alt {
|
||||
return upstream_url;
|
||||
}
|
||||
|
||||
if upstream_url.contains('?') {
|
||||
format!("{upstream_url}&alt=sse")
|
||||
} else {
|
||||
format!("{upstream_url}?alt=sse")
|
||||
}
|
||||
}
|
||||
|
||||
fn custom_path_template_regex() -> &'static Regex {
|
||||
static REGEX: OnceLock<Regex> = OnceLock::new();
|
||||
REGEX.get_or_init(|| {
|
||||
Regex::new(r"\{([A-Za-z_][A-Za-z0-9_]*)\}")
|
||||
.expect("custom_path template regex should compile")
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{build_transport_request_url, TransportRequestUrlParams};
|
||||
use crate::snapshot::{
|
||||
GatewayProviderTransportEndpoint, GatewayProviderTransportKey,
|
||||
GatewayProviderTransportProvider, GatewayProviderTransportSnapshot,
|
||||
};
|
||||
|
||||
fn sample_transport(
|
||||
provider_type: &str,
|
||||
api_format: &str,
|
||||
base_url: &str,
|
||||
custom_path: Option<&str>,
|
||||
) -> GatewayProviderTransportSnapshot {
|
||||
GatewayProviderTransportSnapshot {
|
||||
provider: GatewayProviderTransportProvider {
|
||||
id: "provider-1".to_string(),
|
||||
name: "provider".to_string(),
|
||||
provider_type: provider_type.to_string(),
|
||||
website: None,
|
||||
is_active: true,
|
||||
keep_priority_on_conversion: false,
|
||||
enable_format_conversion: false,
|
||||
concurrent_limit: None,
|
||||
max_retries: None,
|
||||
proxy: None,
|
||||
request_timeout_secs: None,
|
||||
stream_first_byte_timeout_secs: None,
|
||||
config: None,
|
||||
},
|
||||
endpoint: GatewayProviderTransportEndpoint {
|
||||
id: "endpoint-1".to_string(),
|
||||
provider_id: "provider-1".to_string(),
|
||||
api_format: api_format.to_string(),
|
||||
api_family: None,
|
||||
endpoint_kind: None,
|
||||
is_active: true,
|
||||
base_url: base_url.to_string(),
|
||||
header_rules: None,
|
||||
body_rules: None,
|
||||
max_retries: None,
|
||||
custom_path: custom_path.map(ToOwned::to_owned),
|
||||
config: None,
|
||||
format_acceptance_config: None,
|
||||
proxy: None,
|
||||
},
|
||||
key: GatewayProviderTransportKey {
|
||||
id: "key-1".to_string(),
|
||||
provider_id: "provider-1".to_string(),
|
||||
name: "key".to_string(),
|
||||
auth_type: "api_key".to_string(),
|
||||
is_active: true,
|
||||
api_formats: None,
|
||||
allowed_models: None,
|
||||
capabilities: None,
|
||||
rate_multipliers: None,
|
||||
global_priority_by_format: None,
|
||||
expires_at_unix_secs: None,
|
||||
proxy: None,
|
||||
fingerprint: None,
|
||||
decrypted_api_key: "vertex-secret".to_string(),
|
||||
decrypted_auth_config: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uses_vertex_hook_before_custom_path_for_custom_aiplatform_transport() {
|
||||
let transport = sample_transport(
|
||||
"custom",
|
||||
"gemini:cli",
|
||||
"https://aiplatform.googleapis.com",
|
||||
Some("/custom/{model}:{action}"),
|
||||
);
|
||||
|
||||
let url = build_transport_request_url(
|
||||
&transport,
|
||||
TransportRequestUrlParams {
|
||||
provider_api_format: "gemini:cli",
|
||||
mapped_model: Some("gemini-3.1-pro-preview"),
|
||||
upstream_is_stream: true,
|
||||
request_query: Some("foo=bar"),
|
||||
kiro_api_region: None,
|
||||
},
|
||||
)
|
||||
.expect("vertex hook url");
|
||||
|
||||
assert_eq!(
|
||||
url,
|
||||
"https://aiplatform.googleapis.com/v1/publishers/google/models/gemini-3.1-pro-preview:streamGenerateContent?alt=sse&foo=bar&key=vertex-secret"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn expands_custom_path_templates_when_hook_does_not_apply() {
|
||||
let transport = sample_transport(
|
||||
"custom",
|
||||
"gemini:chat",
|
||||
"https://generativelanguage.googleapis.com",
|
||||
Some("/v1beta/models/{model}:{action}"),
|
||||
);
|
||||
|
||||
let url = build_transport_request_url(
|
||||
&transport,
|
||||
TransportRequestUrlParams {
|
||||
provider_api_format: "gemini:chat",
|
||||
mapped_model: Some("gemini-2.5-pro"),
|
||||
upstream_is_stream: false,
|
||||
request_query: Some("key=client-key&foo=bar"),
|
||||
kiro_api_region: None,
|
||||
},
|
||||
)
|
||||
.expect("expanded custom path url");
|
||||
|
||||
assert_eq!(
|
||||
url,
|
||||
"https://generativelanguage.googleapis.com/v1beta/models/gemini-2.5-pro:generateContent?foo=bar"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn keeps_original_custom_path_when_template_params_are_missing() {
|
||||
let transport = sample_transport(
|
||||
"custom",
|
||||
"claude:chat",
|
||||
"https://api.example.com",
|
||||
Some("/v1/messages/{model}"),
|
||||
);
|
||||
|
||||
let url = build_transport_request_url(
|
||||
&transport,
|
||||
TransportRequestUrlParams {
|
||||
provider_api_format: "claude:chat",
|
||||
mapped_model: None,
|
||||
upstream_is_stream: false,
|
||||
request_query: None,
|
||||
kiro_api_region: None,
|
||||
},
|
||||
)
|
||||
.expect("fallback custom path url");
|
||||
|
||||
assert_eq!(url, "https://api.example.com/v1/messages/{model}");
|
||||
}
|
||||
}
|
||||
@@ -1,10 +1,14 @@
|
||||
mod auth;
|
||||
mod context;
|
||||
mod policy;
|
||||
mod url;
|
||||
|
||||
pub use auth::{
|
||||
resolve_local_vertex_api_key_query_auth, VertexApiKeyQueryAuth, VERTEX_API_KEY_QUERY_PARAM,
|
||||
};
|
||||
pub use context::{
|
||||
is_vertex_api_key_transport_context, looks_like_vertex_ai_host, uses_vertex_api_key_query_auth,
|
||||
};
|
||||
pub use policy::{
|
||||
local_vertex_api_key_gemini_transport_unsupported_reason_with_network,
|
||||
supports_local_vertex_api_key_gemini_transport,
|
||||
|
||||
@@ -11,12 +11,7 @@ pub struct VertexApiKeyQueryAuth {
|
||||
pub fn resolve_local_vertex_api_key_query_auth(
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
) -> Option<VertexApiKeyQueryAuth> {
|
||||
if !transport
|
||||
.provider
|
||||
.provider_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case(super::PROVIDER_TYPE)
|
||||
{
|
||||
if !super::is_vertex_api_key_transport_context(transport) {
|
||||
return None;
|
||||
}
|
||||
|
||||
@@ -127,4 +122,15 @@ mod tests {
|
||||
transport.key.decrypted_auth_config = Some("{\"project_id\":\"demo-project\"}".to_string());
|
||||
assert!(resolve_local_vertex_api_key_query_auth(&transport).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn resolves_query_auth_for_custom_aiplatform_transport() {
|
||||
let mut transport = sample_transport();
|
||||
transport.provider.provider_type = "custom".to_string();
|
||||
transport.endpoint.api_format = "gemini:cli".to_string();
|
||||
|
||||
let auth = resolve_local_vertex_api_key_query_auth(&transport)
|
||||
.expect("custom aiplatform transport should resolve");
|
||||
assert_eq!(auth.value, "vertex-secret");
|
||||
}
|
||||
}
|
||||
|
||||
163
crates/aether-provider-transport/src/vertex/context.rs
Normal file
163
crates/aether-provider-transport/src/vertex/context.rs
Normal file
@@ -0,0 +1,163 @@
|
||||
use url::Url;
|
||||
|
||||
use super::super::snapshot::GatewayProviderTransportSnapshot;
|
||||
|
||||
const VERTEX_AI_HOST: &str = "aiplatform.googleapis.com";
|
||||
|
||||
pub fn looks_like_vertex_ai_host(base_url: &str) -> bool {
|
||||
let trimmed = base_url.trim();
|
||||
if trimmed.is_empty() {
|
||||
return false;
|
||||
}
|
||||
|
||||
let Ok(parsed) = Url::parse(trimmed) else {
|
||||
return false;
|
||||
};
|
||||
let Some(host) = parsed
|
||||
.host_str()
|
||||
.map(|value| value.trim().to_ascii_lowercase())
|
||||
else {
|
||||
return false;
|
||||
};
|
||||
|
||||
host == VERTEX_AI_HOST
|
||||
|| host.ends_with(&format!(".{VERTEX_AI_HOST}"))
|
||||
|| host.ends_with(&format!("-{VERTEX_AI_HOST}"))
|
||||
}
|
||||
|
||||
pub fn is_vertex_api_key_transport_context(transport: &GatewayProviderTransportSnapshot) -> bool {
|
||||
if transport
|
||||
.provider
|
||||
.provider_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case(super::PROVIDER_TYPE)
|
||||
{
|
||||
return transport
|
||||
.key
|
||||
.auth_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case("api_key");
|
||||
}
|
||||
|
||||
if !looks_like_vertex_ai_host(&transport.endpoint.base_url) {
|
||||
return false;
|
||||
}
|
||||
|
||||
let endpoint_api_format = transport.endpoint.api_format.trim().to_ascii_lowercase();
|
||||
if !endpoint_api_format.starts_with("gemini:") && !endpoint_api_format.starts_with("claude:") {
|
||||
return false;
|
||||
}
|
||||
|
||||
transport
|
||||
.key
|
||||
.auth_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case("api_key")
|
||||
}
|
||||
|
||||
pub fn uses_vertex_api_key_query_auth(
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
provider_api_format: &str,
|
||||
) -> bool {
|
||||
is_vertex_api_key_transport_context(transport)
|
||||
&& provider_api_format
|
||||
.trim()
|
||||
.to_ascii_lowercase()
|
||||
.starts_with("gemini:")
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{
|
||||
is_vertex_api_key_transport_context, looks_like_vertex_ai_host,
|
||||
uses_vertex_api_key_query_auth,
|
||||
};
|
||||
use crate::snapshot::{
|
||||
GatewayProviderTransportEndpoint, GatewayProviderTransportKey,
|
||||
GatewayProviderTransportProvider, GatewayProviderTransportSnapshot,
|
||||
};
|
||||
|
||||
fn sample_transport() -> GatewayProviderTransportSnapshot {
|
||||
GatewayProviderTransportSnapshot {
|
||||
provider: GatewayProviderTransportProvider {
|
||||
id: "provider-1".to_string(),
|
||||
name: "Vertex".to_string(),
|
||||
provider_type: "custom".to_string(),
|
||||
website: None,
|
||||
is_active: true,
|
||||
keep_priority_on_conversion: false,
|
||||
enable_format_conversion: false,
|
||||
concurrent_limit: None,
|
||||
max_retries: None,
|
||||
proxy: None,
|
||||
request_timeout_secs: None,
|
||||
stream_first_byte_timeout_secs: None,
|
||||
config: None,
|
||||
},
|
||||
endpoint: GatewayProviderTransportEndpoint {
|
||||
id: "endpoint-1".to_string(),
|
||||
provider_id: "provider-1".to_string(),
|
||||
api_format: "gemini:cli".to_string(),
|
||||
api_family: Some("gemini".to_string()),
|
||||
endpoint_kind: Some("cli".to_string()),
|
||||
is_active: true,
|
||||
base_url: "https://aiplatform.googleapis.com".to_string(),
|
||||
header_rules: None,
|
||||
body_rules: None,
|
||||
max_retries: None,
|
||||
custom_path: Some("/v1/publishers/google/models/{model}:{action}".to_string()),
|
||||
config: None,
|
||||
format_acceptance_config: None,
|
||||
proxy: None,
|
||||
},
|
||||
key: GatewayProviderTransportKey {
|
||||
id: "key-1".to_string(),
|
||||
provider_id: "provider-1".to_string(),
|
||||
name: "key".to_string(),
|
||||
auth_type: "api_key".to_string(),
|
||||
is_active: true,
|
||||
api_formats: Some(vec!["gemini:cli".to_string()]),
|
||||
allowed_models: None,
|
||||
capabilities: None,
|
||||
rate_multipliers: None,
|
||||
global_priority_by_format: None,
|
||||
expires_at_unix_secs: None,
|
||||
proxy: None,
|
||||
fingerprint: None,
|
||||
decrypted_api_key: "vertex-secret".to_string(),
|
||||
decrypted_auth_config: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn detects_vertex_host() {
|
||||
assert!(looks_like_vertex_ai_host(
|
||||
"https://aiplatform.googleapis.com"
|
||||
));
|
||||
assert!(looks_like_vertex_ai_host(
|
||||
"https://us-central1-aiplatform.googleapis.com"
|
||||
));
|
||||
assert!(!looks_like_vertex_ai_host("https://example.com"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn infers_vertex_api_key_context_for_custom_aiplatform_transport() {
|
||||
assert!(is_vertex_api_key_transport_context(&sample_transport()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_non_api_key_custom_aiplatform_transport() {
|
||||
let mut transport = sample_transport();
|
||||
transport.key.auth_type = "bearer".to_string();
|
||||
assert!(!is_vertex_api_key_transport_context(&transport));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn detects_vertex_query_auth_usage_for_gemini_formats() {
|
||||
let transport = sample_transport();
|
||||
assert!(uses_vertex_api_key_query_auth(&transport, "gemini:cli"));
|
||||
assert!(uses_vertex_api_key_query_auth(&transport, "gemini:chat"));
|
||||
assert!(!uses_vertex_api_key_query_auth(&transport, "claude:chat"));
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,15 @@ use super::super::{
|
||||
};
|
||||
use super::auth::resolve_local_vertex_api_key_query_auth;
|
||||
|
||||
fn is_vertex_transport_family(transport: &GatewayProviderTransportSnapshot) -> bool {
|
||||
transport
|
||||
.provider
|
||||
.provider_type
|
||||
.trim()
|
||||
.eq_ignore_ascii_case(super::PROVIDER_TYPE)
|
||||
|| super::looks_like_vertex_ai_host(&transport.endpoint.base_url)
|
||||
}
|
||||
|
||||
pub fn local_vertex_api_key_gemini_transport_unsupported_reason_with_network(
|
||||
transport: &GatewayProviderTransportSnapshot,
|
||||
) -> Option<&'static str> {
|
||||
@@ -18,19 +27,21 @@ pub fn local_vertex_api_key_gemini_transport_unsupported_reason_with_network(
|
||||
};
|
||||
}
|
||||
if !transport
|
||||
.provider
|
||||
.provider_type
|
||||
.endpoint
|
||||
.api_format
|
||||
.trim()
|
||||
.eq_ignore_ascii_case(super::PROVIDER_TYPE)
|
||||
{
|
||||
return Some("transport_provider_type_unsupported");
|
||||
}
|
||||
let endpoint_api_format = transport.endpoint.api_format.trim();
|
||||
if !endpoint_api_format.eq_ignore_ascii_case("gemini:chat")
|
||||
&& !endpoint_api_format.eq_ignore_ascii_case("gemini:cli")
|
||||
.eq_ignore_ascii_case("gemini:chat")
|
||||
&& !transport
|
||||
.endpoint
|
||||
.api_format
|
||||
.trim()
|
||||
.eq_ignore_ascii_case("gemini:cli")
|
||||
{
|
||||
return Some("transport_api_format_mismatch");
|
||||
}
|
||||
if !is_vertex_transport_family(transport) {
|
||||
return Some("transport_provider_type_unsupported");
|
||||
}
|
||||
if !header_rules_are_locally_supported(transport.endpoint.header_rules.as_ref()) {
|
||||
return Some("transport_header_rules_unsupported");
|
||||
}
|
||||
@@ -87,18 +98,21 @@ fn supports_local_vertex_api_key_same_format_transport(
|
||||
return false;
|
||||
}
|
||||
if !transport
|
||||
.provider
|
||||
.provider_type
|
||||
.endpoint
|
||||
.api_format
|
||||
.trim()
|
||||
.eq_ignore_ascii_case(super::PROVIDER_TYPE)
|
||||
.eq_ignore_ascii_case(api_formats[0])
|
||||
&& !api_formats.iter().any(|api_format| {
|
||||
transport
|
||||
.endpoint
|
||||
.api_format
|
||||
.trim()
|
||||
.eq_ignore_ascii_case(api_format)
|
||||
})
|
||||
{
|
||||
return false;
|
||||
}
|
||||
let endpoint_api_format = transport.endpoint.api_format.trim();
|
||||
if !api_formats
|
||||
.iter()
|
||||
.any(|api_format| endpoint_api_format.eq_ignore_ascii_case(api_format))
|
||||
{
|
||||
if !super::is_vertex_api_key_transport_context(transport) {
|
||||
return false;
|
||||
}
|
||||
if !header_rules_are_locally_supported(transport.endpoint.header_rules.as_ref())
|
||||
@@ -223,6 +237,14 @@ mod tests {
|
||||
assert!(supports_local_vertex_api_key_gemini_transport(&transport));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn supports_custom_aiplatform_gemini_cli_subset() {
|
||||
let mut transport = sample_transport();
|
||||
transport.provider.provider_type = "custom".to_string();
|
||||
transport.endpoint.api_format = "gemini:cli".to_string();
|
||||
assert!(supports_local_vertex_api_key_gemini_transport(&transport));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_vertex_service_account_subset() {
|
||||
let mut transport = sample_transport();
|
||||
|
||||
@@ -17,6 +17,9 @@ pub struct SchedulerRequestCandidateReportContext {
|
||||
pub key_id: Option<String>,
|
||||
pub client_api_format: Option<String>,
|
||||
pub provider_api_format: Option<String>,
|
||||
pub upstream_url: Option<String>,
|
||||
pub mapped_model: Option<String>,
|
||||
pub key_name: Option<String>,
|
||||
pub proxy: Option<Value>,
|
||||
}
|
||||
|
||||
@@ -103,6 +106,9 @@ pub fn parse_request_candidate_report_context(
|
||||
key_id: string_field(report_context, "key_id"),
|
||||
client_api_format: string_field(report_context, "client_api_format"),
|
||||
provider_api_format: string_field(report_context, "provider_api_format"),
|
||||
upstream_url: string_field(report_context, "upstream_url"),
|
||||
mapped_model: string_field(report_context, "mapped_model"),
|
||||
key_name: string_field(report_context, "key_name"),
|
||||
proxy: report_context
|
||||
.get("proxy")
|
||||
.cloned()
|
||||
@@ -129,11 +135,20 @@ pub fn resolve_report_request_candidate_slot(
|
||||
key_id,
|
||||
client_api_format,
|
||||
provider_api_format,
|
||||
upstream_url,
|
||||
mapped_model,
|
||||
key_name,
|
||||
proxy,
|
||||
} = metadata;
|
||||
let request_id = request_id?;
|
||||
let synthesized_extra_data =
|
||||
build_report_candidate_extra_data(client_api_format, provider_api_format, proxy);
|
||||
let synthesized_extra_data = build_report_candidate_extra_data(
|
||||
client_api_format,
|
||||
provider_api_format,
|
||||
upstream_url,
|
||||
mapped_model,
|
||||
key_name,
|
||||
proxy,
|
||||
);
|
||||
let created_at_unix_ms = matched_candidate
|
||||
.as_ref()
|
||||
.map(|candidate| candidate.created_at_unix_ms)
|
||||
@@ -489,9 +504,12 @@ fn next_candidate_index(candidates: &[StoredRequestCandidate]) -> u32 {
|
||||
fn build_report_candidate_extra_data(
|
||||
client_api_format: Option<String>,
|
||||
provider_api_format: Option<String>,
|
||||
upstream_url: Option<String>,
|
||||
mapped_model: Option<String>,
|
||||
key_name: Option<String>,
|
||||
proxy: Option<Value>,
|
||||
) -> Option<Value> {
|
||||
let mut extra_data = Map::with_capacity(5);
|
||||
let mut extra_data = Map::with_capacity(8);
|
||||
extra_data.insert("gateway_execution_runtime".to_string(), Value::Bool(true));
|
||||
extra_data.insert("phase".to_string(), Value::String("3c_trial".to_string()));
|
||||
if let Some(client_api_format) = client_api_format {
|
||||
@@ -506,6 +524,15 @@ fn build_report_candidate_extra_data(
|
||||
Value::String(provider_api_format),
|
||||
);
|
||||
}
|
||||
if let Some(upstream_url) = upstream_url {
|
||||
extra_data.insert("upstream_url".to_string(), Value::String(upstream_url));
|
||||
}
|
||||
if let Some(mapped_model) = mapped_model {
|
||||
extra_data.insert("mapped_model".to_string(), Value::String(mapped_model));
|
||||
}
|
||||
if let Some(key_name) = key_name {
|
||||
extra_data.insert("key_name".to_string(), Value::String(key_name));
|
||||
}
|
||||
if let Some(proxy) = proxy {
|
||||
extra_data.insert("proxy".to_string(), proxy);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user