Merge branch 'pr-562' into review/pr-562

This commit is contained in:
fawney19
2026-05-25 23:10:29 +08:00
23 changed files with 1715 additions and 57 deletions
@@ -1441,7 +1441,7 @@ async fn gateway_executes_antigravity_gemini_cli_stream_via_local_decision_gate_
false,
Some(serde_json::json!(["gemini", "antigravity"])),
Some(serde_json::json!(["gemini:generate_content"])),
Some(serde_json::json!(["gemini-cli"])),
Some(serde_json::json!(["gemini-cli", "gemini-3.1-flash-lite"])),
api_key_id.to_string(),
Some("default".to_string()),
true,
@@ -1452,12 +1452,23 @@ async fn gateway_executes_antigravity_gemini_cli_stream_via_local_decision_gate_
Some(4_102_444_800),
Some(serde_json::json!(["gemini", "antigravity"])),
Some(serde_json::json!(["gemini:generate_content"])),
Some(serde_json::json!(["gemini-cli"])),
Some(serde_json::json!(["gemini-cli", "gemini-3.1-flash-lite"])),
)
.expect("auth snapshot should build")
}
fn sample_candidate_row() -> StoredMinimalCandidateSelectionRow {
sample_candidate_row_for("gemini-cli", "1")
}
fn sample_native_antigravity_candidate_row() -> StoredMinimalCandidateSelectionRow {
sample_candidate_row_for("gemini-3.1-flash-lite", "native-1")
}
fn sample_candidate_row_for(
global_model_name: &str,
row_suffix: &str,
) -> StoredMinimalCandidateSelectionRow {
StoredMinimalCandidateSelectionRow {
provider_id: "provider-antigravity-cli-oauth-stream-local-1".to_string(),
provider_name: "antigravity".to_string(),
@@ -1478,9 +1489,11 @@ async fn gateway_executes_antigravity_gemini_cli_stream_via_local_decision_gate_
key_capabilities: None,
key_internal_priority: 5,
key_global_priority_by_format: Some(serde_json::json!({"gemini:generate_content": 1})),
model_id: "model-antigravity-cli-oauth-stream-local-1".to_string(),
global_model_id: "global-model-antigravity-cli-oauth-stream-local-1".to_string(),
global_model_name: "gemini-cli".to_string(),
model_id: format!("model-antigravity-cli-oauth-stream-local-{row_suffix}"),
global_model_id: format!(
"global-model-antigravity-cli-oauth-stream-local-{row_suffix}"
),
global_model_name: global_model_name.to_string(),
global_model_mappings: None,
global_model_supports_streaming: Some(true),
model_provider_model_name: "claude-sonnet-4-5".to_string(),
@@ -1800,6 +1813,7 @@ async fn gateway_executes_antigravity_gemini_cli_stream_via_local_decision_gate_
let candidate_selection_repository =
Arc::new(InMemoryMinimalCandidateSelectionReadRepository::seed(vec![
sample_candidate_row(),
sample_native_antigravity_candidate_row(),
]));
let request_candidate_repository = Arc::new(InMemoryRequestCandidateRepository::default());
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
@@ -1818,17 +1832,26 @@ async fn gateway_executes_antigravity_gemini_cli_stream_via_local_decision_gate_
.with_token_url_for_tests("antigravity", format!("{refresh_url}/oauth/token")),
),
]);
let gateway_state = build_state_with_execution_runtime_override(execution_runtime_url.clone())
.with_data_state_for_tests(
let data_state =
crate::data::GatewayDataState::with_auth_candidate_selection_provider_catalog_and_request_candidate_repository_for_tests(
auth_repository,
candidate_selection_repository,
provider_catalog_repository,
Arc::clone(&request_candidate_repository),
DEVELOPMENT_ENCRYPTION_KEY,
),
)
.with_oauth_refresh_coordinator_for_tests(oauth_refresh);
)
.with_system_config_values_for_tests([(
crate::constants::ANTIGRAVITY_BEARER_BRIDGE_CONFIG_KEY.to_string(),
json!({
"enabled": true,
"auth_user_id": "user-antigravity-cli-oauth-stream-local-1",
"auth_api_key_id": "api-key-antigravity-cli-oauth-stream-local-1",
"allow_unverified_google_bearer": true
}),
)]);
let gateway_state = build_state_with_execution_runtime_override(execution_runtime_url.clone())
.with_data_state_for_tests(data_state)
.with_oauth_refresh_coordinator_for_tests(oauth_refresh);
let gateway = build_router_with_state(gateway_state);
let (gateway_url, gateway_handle) = start_server(gateway).await;
@@ -1950,6 +1973,380 @@ async fn gateway_executes_antigravity_gemini_cli_stream_via_local_decision_gate_
assert_eq!(stored_candidates.len(), 1);
assert_eq!(stored_candidates[0].status, RequestCandidateStatus::Success);
*seen_execution_runtime.lock().expect("mutex should lock") = None;
let inbound_response = reqwest::Client::new()
.post(format!(
"{gateway_url}/v1internal:streamGenerateContent?alt=sse"
))
.header(http::header::CONTENT_TYPE, "application/json")
.header("authorization", "Bearer google-antigravity-access-token")
.header("x-api-key", client_api_key)
.header("user-agent", "antigravity/cli/1.0.2 linux/arm64")
.header(
TRACE_ID_HEADER,
"trace-antigravity-v1internal-inbound-stream-456",
)
.json(&json!({
"project": "client-side-project-should-not-leak",
"requestId": "client-v1internal-request-456",
"model": "gemini-cli",
"userAgent": "antigravity",
"requestType": "checkpoint",
"request": {
"contents": [{
"role": "user",
"parts": [{"text": "checkpoint context"}]
}],
"generationConfig": {
"temperature": 0.4,
"thinkingConfig": {
"includeThoughts": true
}
},
"toolConfig": {
"functionCallingConfig": {
"mode": "NONE"
}
}
}
}))
.send()
.await
.expect("inbound antigravity request should succeed");
let inbound_status = inbound_response.status();
let inbound_miss_reason = inbound_response
.headers()
.get(crate::constants::LOCAL_EXECUTION_RUNTIME_MISS_REASON_HEADER)
.and_then(|value| value.to_str().ok())
.unwrap_or("-")
.to_string();
let inbound_response_body = inbound_response.text().await.expect("body should read");
assert_eq!(
inbound_status,
StatusCode::OK,
"unexpected inbound antigravity response body: {inbound_response_body}; miss_reason={inbound_miss_reason}"
);
let inbound_response_text = strip_sse_keepalive_comments(&inbound_response_body);
let inbound_payload = inbound_response_text
.trim()
.strip_prefix("data: ")
.expect("response should start with sse data prefix");
let inbound_response_json: serde_json::Value =
serde_json::from_str(inbound_payload).expect("stream payload should parse");
assert_eq!(
inbound_response_json["responseId"],
"resp_antigravity_cli_local_stream_123"
);
assert_eq!(
inbound_response_json["response"]["candidates"][0]["content"]["parts"][0]["text"],
"Hello Antigravity Stream"
);
let seen_inbound_execution_runtime_request = seen_execution_runtime
.lock()
.expect("mutex should lock")
.clone()
.expect("inbound execution runtime stream should be captured");
assert_eq!(
seen_inbound_execution_runtime_request.trace_id,
"trace-antigravity-v1internal-inbound-stream-456"
);
assert_eq!(
seen_inbound_execution_runtime_request.url,
"https://antigravity.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert_eq!(
seen_inbound_execution_runtime_request.authorization,
"Bearer refreshed-antigravity-cli-stream-access-token"
);
assert_eq!(
seen_inbound_execution_runtime_request.project,
"project-antigravity-stream-local-1"
);
assert_eq!(
seen_inbound_execution_runtime_request.request_id,
"client-v1internal-request-456"
);
assert_eq!(
seen_inbound_execution_runtime_request.model,
"claude-sonnet-4-5"
);
assert_eq!(
seen_inbound_execution_runtime_request.user_agent,
"antigravity"
);
assert_eq!(
seen_inbound_execution_runtime_request.request_type,
"checkpoint"
);
assert_eq!(seen_inbound_execution_runtime_request.contents_len, 1);
assert!((seen_inbound_execution_runtime_request.exact_temperature - 0.4).abs() < f64::EPSILON);
assert!(!seen_inbound_execution_runtime_request.request_has_model);
*seen_execution_runtime.lock().expect("mutex should lock") = None;
let bearer_only_response = reqwest::Client::new()
.post(format!(
"{gateway_url}/v1internal:streamGenerateContent?alt=sse"
))
.header(http::header::CONTENT_TYPE, "application/json")
.header("authorization", "Bearer google-antigravity-access-token")
.header("user-agent", "antigravity/cli/1.0.2 linux/arm64")
.header(
TRACE_ID_HEADER,
"trace-antigravity-v1internal-bearer-only-stream-789",
)
.json(&json!({
"project": "client-side-project-should-not-leak",
"requestId": "client-v1internal-request-789",
"model": "gemini-cli",
"userAgent": "antigravity",
"requestType": "agent",
"request": {
"contents": [{
"role": "user",
"parts": [{"text": "bearer-only request"}]
}],
"generationConfig": {
"temperature": 0.5,
"thinkingConfig": {
"includeThoughts": true
}
},
"toolConfig": {
"functionCallingConfig": {
"mode": "NONE"
}
}
}
}))
.send()
.await
.expect("bearer-only antigravity request should succeed");
let bearer_only_status = bearer_only_response.status();
let bearer_only_miss_reason = bearer_only_response
.headers()
.get(crate::constants::LOCAL_EXECUTION_RUNTIME_MISS_REASON_HEADER)
.and_then(|value| value.to_str().ok())
.unwrap_or("-")
.to_string();
let bearer_only_response_body = bearer_only_response.text().await.expect("body should read");
assert_eq!(
bearer_only_status,
StatusCode::OK,
"unexpected bearer-only antigravity response body: {bearer_only_response_body}; miss_reason={bearer_only_miss_reason}"
);
let bearer_only_response_text = strip_sse_keepalive_comments(&bearer_only_response_body);
let bearer_only_payload = bearer_only_response_text
.trim()
.strip_prefix("data: ")
.expect("response should start with sse data prefix");
let bearer_only_response_json: serde_json::Value =
serde_json::from_str(bearer_only_payload).expect("stream payload should parse");
assert_eq!(
bearer_only_response_json["responseId"],
"resp_antigravity_cli_local_stream_123"
);
assert_eq!(
bearer_only_response_json["response"]["candidates"][0]["content"]["parts"][0]["text"],
"Hello Antigravity Stream"
);
let seen_bearer_only_execution_runtime_request = seen_execution_runtime
.lock()
.expect("mutex should lock")
.clone()
.expect("bearer-only inbound execution runtime stream should be captured");
assert_eq!(
seen_bearer_only_execution_runtime_request.trace_id,
"trace-antigravity-v1internal-bearer-only-stream-789"
);
assert_eq!(
seen_bearer_only_execution_runtime_request.url,
"https://antigravity.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert_eq!(
seen_bearer_only_execution_runtime_request.authorization,
"Bearer refreshed-antigravity-cli-stream-access-token"
);
assert_eq!(
seen_bearer_only_execution_runtime_request.request_id,
"client-v1internal-request-789"
);
assert_eq!(
seen_bearer_only_execution_runtime_request.request_type,
"agent"
);
assert_eq!(seen_bearer_only_execution_runtime_request.contents_len, 1);
assert!(
(seen_bearer_only_execution_runtime_request.exact_temperature - 0.5).abs() < f64::EPSILON
);
assert!(!seen_bearer_only_execution_runtime_request.request_has_model);
*seen_execution_runtime.lock().expect("mutex should lock") = None;
let native_model_response = reqwest::Client::new()
.post(format!(
"{gateway_url}/v1internal:streamGenerateContent?alt=sse"
))
.header(http::header::CONTENT_TYPE, "application/json")
.header("authorization", "Bearer google-antigravity-access-token")
.header("user-agent", "antigravity/cli/1.0.2 linux/arm64")
.header(
TRACE_ID_HEADER,
"trace-antigravity-v1internal-native-model-stream-790",
)
.json(&json!({
"project": "client-side-project-should-not-leak",
"requestId": "client-v1internal-request-790",
"model": "gemini-3.1-flash-lite",
"userAgent": "antigravity",
"requestType": "agent",
"request": {
"contents": [{
"role": "user",
"parts": [{"text": "native antigravity model request"}]
}],
"generationConfig": {
"temperature": 0.6,
"thinkingConfig": {
"includeThoughts": true
}
},
"toolConfig": {
"functionCallingConfig": {
"mode": "NONE"
}
}
}
}))
.send()
.await
.expect("native-model antigravity request should succeed");
let native_model_status = native_model_response.status();
let native_model_miss_reason = native_model_response
.headers()
.get(crate::constants::LOCAL_EXECUTION_RUNTIME_MISS_REASON_HEADER)
.and_then(|value| value.to_str().ok())
.unwrap_or("-")
.to_string();
let native_model_response_body = native_model_response
.text()
.await
.expect("body should read");
assert_eq!(
native_model_status,
StatusCode::OK,
"unexpected native-model antigravity response body: {native_model_response_body}; miss_reason={native_model_miss_reason}"
);
let native_model_response_text = strip_sse_keepalive_comments(&native_model_response_body);
let native_model_payload = native_model_response_text
.trim()
.strip_prefix("data: ")
.expect("response should start with sse data prefix");
let native_model_response_json: serde_json::Value =
serde_json::from_str(native_model_payload).expect("stream payload should parse");
assert_eq!(
native_model_response_json["responseId"],
"resp_antigravity_cli_local_stream_123"
);
assert_eq!(
native_model_response_json["response"]["candidates"][0]["content"]["parts"][0]["text"],
"Hello Antigravity Stream"
);
let seen_native_model_execution_runtime_request = seen_execution_runtime
.lock()
.expect("mutex should lock")
.clone()
.expect("native-model inbound execution runtime stream should be captured");
assert_eq!(
seen_native_model_execution_runtime_request.trace_id,
"trace-antigravity-v1internal-native-model-stream-790"
);
assert_eq!(
seen_native_model_execution_runtime_request.url,
"https://antigravity.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert_eq!(
seen_native_model_execution_runtime_request.authorization,
"Bearer refreshed-antigravity-cli-stream-access-token"
);
assert_eq!(
seen_native_model_execution_runtime_request.model,
"claude-sonnet-4-5"
);
assert_eq!(
seen_native_model_execution_runtime_request.request_id,
"client-v1internal-request-790"
);
assert!(
(seen_native_model_execution_runtime_request.exact_temperature - 0.6).abs() < f64::EPSILON
);
assert!(!seen_native_model_execution_runtime_request.request_has_model);
if std::env::var("AETHER_REAL_AGY_CLI_SMOKE").ok().as_deref() == Some("1") {
*seen_execution_runtime.lock().expect("mutex should lock") = None;
let log_path = std::env::var("AETHER_REAL_AGY_CLI_LOG")
.unwrap_or_else(|_| "/tmp/aether-real-agy-cli-smoke.log".to_string());
let workdir = std::env::var("AETHER_REAL_AGY_CLI_WORKDIR")
.unwrap_or_else(|_| "/tmp/aether-real-agy-cli-work".to_string());
std::fs::create_dir_all(&workdir).expect("agy smoke workdir should create");
let gateway_url_for_agy = gateway_url.clone();
let log_path_for_agy = log_path.clone();
let workdir_for_agy = workdir.clone();
let output = tokio::task::spawn_blocking(move || {
std::process::Command::new("agy")
.arg("--log-file")
.arg(&log_path_for_agy)
.arg("-p")
.arg("Reply with AETHER_CLOSED_LOOP_OK only.")
.arg("--print-timeout")
.arg("45s")
.env("AGY_CLI_DISABLE_AUTO_UPDATE", "true")
.env("CLOUD_CODE_URL", &gateway_url_for_agy)
.current_dir(&workdir_for_agy)
.output()
})
.await
.expect("agy smoke blocking task should join")
.expect("agy smoke process should spawn");
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
let agy_log = std::fs::read_to_string(&log_path).unwrap_or_default();
let seen_agy_execution_runtime_snapshot = seen_execution_runtime
.lock()
.expect("mutex should lock")
.clone();
assert!(
output.status.success(),
"agy smoke failed: status={:?}\nseen_execution_runtime={seen_agy_execution_runtime_snapshot:?}\nstdout={stdout}\nstderr={stderr}\nlog={agy_log}",
output.status
);
assert!(
stdout.contains("Hello Antigravity Stream")
|| stdout.contains("AETHER_CLOSED_LOOP_OK"),
"agy smoke stdout did not contain the local runtime response: stdout={stdout}\nstderr={stderr}\nlog={agy_log}"
);
let seen_agy_execution_runtime_request = seen_agy_execution_runtime_snapshot
.expect("real agy smoke should reach execution runtime");
assert_eq!(
seen_agy_execution_runtime_request.url,
"https://antigravity.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
}
let inbound_stored_candidates = request_candidate_repository
.list_by_request_id("trace-antigravity-v1internal-inbound-stream-456")
.await
.expect("inbound request candidate trace should read");
assert_eq!(inbound_stored_candidates.len(), 1);
assert_eq!(
inbound_stored_candidates[0].status,
RequestCandidateStatus::Success
);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
assert!(
!*seen_report.lock().expect("mutex should lock"),
@@ -700,6 +700,174 @@ async fn gateway_rejects_invalid_claude_count_tokens_payload_without_hitting_fal
fallback_probe_handle.abort();
}
#[tokio::test]
async fn gateway_handles_antigravity_v1internal_control_plane_without_proxying() {
let fallback_probe_hits = Arc::new(Mutex::new(0usize));
let fallback_probe_hits_clone = Arc::clone(&fallback_probe_hits);
let fallback_probe = Router::new().route(
"/{*path}",
any(move |_request: Request| {
let fallback_probe_hits_inner = Arc::clone(&fallback_probe_hits_clone);
async move {
*fallback_probe_hits_inner.lock().expect("mutex should lock") += 1;
(StatusCode::OK, Json(json!({"proxied": true}))).into_response()
}
}),
);
let (_unused_fallback_probe_url, fallback_probe_handle) = start_server(fallback_probe).await;
let gateway = build_router_with_state(AppState::new().expect("gateway should build"));
let (gateway_url, gateway_handle) = start_server(gateway).await;
let client = reqwest::Client::new();
let user_settings = json!({
"preferredModelId": "gemini-3.1-flash-lite",
"theme": "dark"
});
let requests = vec![
(
"/v1internal:loadCodeAssist",
json!({"metadata": {"ideType": "ANTIGRAVITY_CLI"}}),
),
(
"/v1internal:fetchAvailableModels",
json!({"project": "aether-antigravity-local"}),
),
(
"/v1internal:fetchUserInfo",
json!({"project": "aether-antigravity-local"}),
),
(
"/v1internal:fetchAdminControls",
json!({"project": "aether-antigravity-local"}),
),
("/v1internal:listExperiments", json!({})),
(
"/v1internal:recordCodeAssistMetrics",
json!({
"project": "aether-antigravity-local",
"requestId": "opaque-request-id",
"metrics": []
}),
),
(
"/v1internal:setUserSettings",
json!({"userSettings": user_settings.clone()}),
),
];
for (path, request_body) in requests {
let response = client
.post(format!("{gateway_url}{path}"))
.header("authorization", "Bearer ant-access-token")
.header("user-agent", "antigravity/cli/1.0.2 linux/arm64")
.json(&request_body)
.send()
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK, "path {path}");
assert_eq!(
response
.headers()
.get(EXECUTION_PATH_HEADER)
.and_then(|value| value.to_str().ok()),
Some(EXECUTION_PATH_LOCAL_AI_PUBLIC),
"path {path}"
);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
match path {
"/v1internal:loadCodeAssist" => {
assert_eq!(
payload["cloudaicompanionProject"],
"aether-antigravity-local"
);
assert_eq!(payload["currentTier"]["id"], "free-tier");
assert_eq!(payload["currentTier"]["name"], "Antigravity");
assert_eq!(payload["paidTier"]["id"], "g1-pro-tier");
assert_eq!(payload["gcpManaged"], false);
assert_eq!(payload["allowedTiers"][0]["id"], "free-tier");
assert_eq!(payload["allowedTiers"][0]["isDefault"], true);
assert_eq!(payload["allowedTiers"][1]["id"], "standard-tier");
assert_eq!(
payload["upgradeSubscriptionUri"],
"https://codeassist.google.com/upgrade"
);
}
"/v1internal:fetchAvailableModels" => {
assert_eq!(payload["defaultAgentModelId"], "gemini-3.1-flash-lite");
assert_eq!(
payload["tieredModelIds"]["flash"],
json!(["gemini-3-flash-agent"])
);
assert_eq!(
payload["models"]["gemini-3.5-flash-low"]["displayName"],
"Gemini 3.5 Flash Low"
);
assert_eq!(
payload["models"]["gemini-3.5-flash-low"]["apiProvider"],
"API_PROVIDER_GOOGLE_GEMINI"
);
assert_eq!(
payload["models"]["gemini-2.5-flash-lite"]["model"],
"MODEL_GOOGLE_GEMINI_2_5_FLASH_LITE"
);
assert_eq!(
payload["agentModelSorts"][0]["groups"][0]["modelIds"],
json!([
"gemini-3.1-flash-lite",
"gemini-3-flash-agent",
"gemini-3.1-pro-low",
"gemini-3.5-flash-low"
])
);
assert_eq!(payload["deprecatedModelIds"], json!({}));
assert_eq!(payload["commandModelIds"], json!(["gemini-3-flash"]));
assert_eq!(
payload["imageGenerationModelIds"],
json!(["gemini-3.1-flash-image"])
);
assert_eq!(payload["mqueryModelIds"], json!(["gemini-3.1-flash-lite"]));
assert_eq!(
payload["webSearchModelIds"],
json!(["gemini-3.1-flash-lite"])
);
assert_eq!(
payload["commitMessageModelIds"],
json!(["gemini-3.1-flash-lite"])
);
}
"/v1internal:fetchUserInfo" => {
assert_eq!(payload["regionCode"], "US");
assert_eq!(
payload["userSettings"]["preferredModelId"],
"gemini-3.1-flash-lite"
);
}
"/v1internal:fetchAdminControls" => {
assert_eq!(payload, json!({}));
}
"/v1internal:listExperiments" => {
assert_eq!(payload["experimentIds"], json!([]));
assert_eq!(payload["flags"], json!([]));
}
"/v1internal:recordCodeAssistMetrics" => {
assert_eq!(payload, json!({}));
}
"/v1internal:setUserSettings" => {
assert_eq!(payload["userSettings"], user_settings);
}
other => panic!("unexpected path {other}"),
}
}
assert_eq!(*fallback_probe_hits.lock().expect("mutex should lock"), 0);
gateway_handle.abort();
fallback_probe_handle.abort();
}
#[tokio::test]
async fn gateway_does_not_locally_reject_image_model_name_on_chat_completions() {
let fallback_probe_hits = Arc::new(Mutex::new(0usize));
@@ -168,6 +168,15 @@ async fn gateway_exposes_frontdoor_manifest_without_proxying_upstream() {
assert!(owned_routes
.iter()
.any(|value| value == "/v1beta/files/{path...}"));
assert!(owned_routes
.iter()
.any(|value| value == "/v1internal:loadCodeAssist"));
assert!(owned_routes
.iter()
.any(|value| value == "/v1internal:fetchAvailableModels"));
assert!(owned_routes
.iter()
.any(|value| value == "/v1internal:streamGenerateContent"));
assert_eq!(
payload["rust_frontdoor"]["internal_gateway"]["status"],
"rust_native_control_plane"