mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-06 01:17:46 +08:00
merge(main): resolve subscription usage policy conflicts
This commit is contained in:
@@ -9,6 +9,7 @@ use crate::ai_serving::planner::candidate_preparation::{
|
||||
};
|
||||
use crate::ai_serving::planner::spec_metadata::local_openai_image_spec_metadata;
|
||||
use crate::ai_serving::pure::normalize_openai_image_request_with_options;
|
||||
use crate::ai_serving::transport::antigravity::is_antigravity_provider_transport;
|
||||
use crate::ai_serving::transport::{
|
||||
build_grok_browser_headers, build_grok_upstream_url, build_openai_image_headers,
|
||||
build_openai_image_upstream_url, build_standard_provider_request_headers,
|
||||
@@ -338,6 +339,25 @@ async fn resolve_local_openai_image_to_gemini_candidate_payload_parts(
|
||||
let candidate = &attempt.eligible.candidate;
|
||||
let transport = &attempt.eligible.transport;
|
||||
let provider_api_format = "gemini:generate_content";
|
||||
|
||||
// The gemini:generate_content URL hook rewrites an Antigravity endpoint to
|
||||
// /v1internal:, and this image path has no v1internal envelope to match it.
|
||||
// Skip the candidate instead of posting a bare Gemini body that upstream
|
||||
// would only reject.
|
||||
if is_antigravity_provider_transport(transport) {
|
||||
mark_skipped_local_openai_image_candidate(
|
||||
state,
|
||||
input,
|
||||
trace_id,
|
||||
candidate,
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"transport_unsupported",
|
||||
)
|
||||
.await;
|
||||
return None;
|
||||
}
|
||||
|
||||
let effective_headers = input.effective_headers(&parts.headers);
|
||||
|
||||
let prepared_candidate = match prepare_header_authenticated_candidate(
|
||||
|
||||
@@ -505,7 +505,7 @@ fn projects_uuid_prompt_cache_identity_into_missing_session_headers() {
|
||||
assert_eq!(headers.get("x-client-request-id"), None);
|
||||
assert_eq!(
|
||||
headers.get("user-agent"),
|
||||
Some(&"codex_cli_rs/0.144.1".to_string())
|
||||
Some(&"codex_cli_rs/0.153.3".to_string())
|
||||
);
|
||||
assert_eq!(headers.get("originator"), Some(&"codex_cli_rs".to_string()));
|
||||
assert!(!headers.contains_key("version"));
|
||||
@@ -615,7 +615,7 @@ fn injects_only_codex_client_headers_for_images_requests() {
|
||||
);
|
||||
assert_eq!(
|
||||
headers.get("user-agent"),
|
||||
Some(&"codex_cli_rs/0.144.1".to_string())
|
||||
Some(&"codex_cli_rs/0.153.3".to_string())
|
||||
);
|
||||
assert_eq!(headers.get("originator"), Some(&"codex_cli_rs".to_string()));
|
||||
assert!(!headers.contains_key("version"));
|
||||
@@ -699,7 +699,7 @@ fn preserves_client_context_headers_and_enforces_codex_provider_identity() {
|
||||
);
|
||||
assert_eq!(
|
||||
headers.get("user-agent"),
|
||||
Some(&"codex_cli_rs/0.144.1".to_string())
|
||||
Some(&"codex_cli_rs/0.153.3".to_string())
|
||||
);
|
||||
assert_eq!(headers.get("originator"), Some(&"codex_cli_rs".to_string()));
|
||||
assert_eq!(
|
||||
@@ -763,7 +763,7 @@ fn compact_projects_uuid_prompt_cache_identity_into_session_headers() {
|
||||
assert_eq!(headers.get("x-client-request-id"), None);
|
||||
assert_eq!(
|
||||
headers.get("user-agent"),
|
||||
Some(&"codex_cli_rs/0.144.1".to_string())
|
||||
Some(&"codex_cli_rs/0.153.3".to_string())
|
||||
);
|
||||
assert_eq!(headers.get("originator"), Some(&"codex_cli_rs".to_string()));
|
||||
assert!(!headers.contains_key("version"));
|
||||
|
||||
@@ -4,6 +4,10 @@ use std::sync::Arc;
|
||||
use aether_contracts::ResolvedTransportProfile;
|
||||
use serde_json::Value;
|
||||
|
||||
use crate::ai_serving::planner::antigravity::{
|
||||
build_antigravity_v1internal_provider_request, AntigravityV1InternalRequestError,
|
||||
AntigravityV1InternalRequestInput, ANTIGRAVITY_V1INTERNAL_ENVELOPE_NAME,
|
||||
};
|
||||
use crate::ai_serving::planner::candidate_preparation::{
|
||||
prepare_header_authenticated_candidate, prepare_header_authenticated_candidate_from_auth,
|
||||
OauthPreparationContext,
|
||||
@@ -26,6 +30,7 @@ use crate::ai_serving::planner::standard::{
|
||||
openai_provider_request_contract_failure_extra_data, openai_responses_reasoning_replay_policy,
|
||||
request_body_build_failure_extra_data, request_conversion_failure_extra_data,
|
||||
};
|
||||
use crate::ai_serving::transport::antigravity::is_antigravity_provider_transport;
|
||||
use crate::ai_serving::transport::kiro::{
|
||||
build_kiro_provider_headers, build_kiro_provider_request_body,
|
||||
is_kiro_claude_messages_transport, KiroProviderHeadersInput, KiroRequestAuth,
|
||||
@@ -837,6 +842,29 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
|
||||
.await);
|
||||
}
|
||||
|
||||
if normalized_provider_api_format == "gemini:generate_content"
|
||||
&& is_antigravity_provider_transport(transport)
|
||||
{
|
||||
return Ok(build_antigravity_cross_format_payload_parts(
|
||||
state,
|
||||
parts,
|
||||
trace_id,
|
||||
body_json,
|
||||
input,
|
||||
attempt,
|
||||
transport,
|
||||
spec_metadata.api_format,
|
||||
provider_api_format,
|
||||
prepared_candidate.mapped_model,
|
||||
prepared_candidate.auth_header,
|
||||
prepared_candidate.auth_value,
|
||||
provider_request_body,
|
||||
upstream_is_stream,
|
||||
redaction.redacted,
|
||||
)
|
||||
.await);
|
||||
}
|
||||
|
||||
if normalized_provider_api_format == "gemini:generate_content"
|
||||
&& is_gemini_cli_provider_transport(transport)
|
||||
{
|
||||
@@ -963,6 +991,145 @@ fn apply_transport_request_body_semantics(
|
||||
)
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn build_antigravity_cross_format_payload_parts(
|
||||
state: &AppState,
|
||||
parts: &http::request::Parts,
|
||||
trace_id: &str,
|
||||
original_body_json: &serde_json::Value,
|
||||
input: &LocalStandardDecisionInput,
|
||||
attempt: &LocalStandardCandidateAttempt,
|
||||
transport: &Arc<GatewayProviderTransportSnapshot>,
|
||||
client_api_format: &str,
|
||||
provider_api_format: &str,
|
||||
mapped_model: String,
|
||||
auth_header: String,
|
||||
auth_value: String,
|
||||
gemini_request_body: Value,
|
||||
upstream_is_stream: bool,
|
||||
request_redacted: bool,
|
||||
) -> Option<LocalStandardCandidatePayloadParts> {
|
||||
let candidate = &attempt.eligible.candidate;
|
||||
let effective_headers = input.effective_headers(&parts.headers);
|
||||
let resolved =
|
||||
match build_antigravity_v1internal_provider_request(AntigravityV1InternalRequestInput {
|
||||
state,
|
||||
parts,
|
||||
transport,
|
||||
trace_id,
|
||||
mapped_model: &mapped_model,
|
||||
provider_api_format,
|
||||
auth_header: &auth_header,
|
||||
auth_value: &auth_value,
|
||||
request_headers: effective_headers,
|
||||
original_request_body: original_body_json,
|
||||
gemini_request_body: &gemini_request_body,
|
||||
upstream_is_stream,
|
||||
same_format: false,
|
||||
})
|
||||
.await
|
||||
{
|
||||
Ok(resolved) => resolved,
|
||||
Err(AntigravityV1InternalRequestError::TransportUnsupported) => {
|
||||
mark_skipped_local_standard_candidate(
|
||||
state,
|
||||
input,
|
||||
trace_id,
|
||||
candidate,
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"transport_unsupported",
|
||||
)
|
||||
.await;
|
||||
return None;
|
||||
}
|
||||
Err(AntigravityV1InternalRequestError::EnvelopeUnsupported) => {
|
||||
mark_skipped_local_standard_candidate_with_extra_data(
|
||||
state,
|
||||
input,
|
||||
trace_id,
|
||||
candidate,
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"provider_request_body_build_failed",
|
||||
request_body_build_failure_extra_data(
|
||||
original_body_json,
|
||||
client_api_format,
|
||||
provider_api_format,
|
||||
),
|
||||
)
|
||||
.await;
|
||||
return None;
|
||||
}
|
||||
Err(AntigravityV1InternalRequestError::UpstreamUrlUnavailable) => {
|
||||
mark_skipped_local_standard_candidate_with_failure_diagnostic(
|
||||
state,
|
||||
input,
|
||||
trace_id,
|
||||
candidate,
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"upstream_url_missing",
|
||||
CandidateFailureDiagnostic::upstream_url_missing(
|
||||
client_api_format,
|
||||
provider_api_format,
|
||||
"standard_family_antigravity_url",
|
||||
),
|
||||
)
|
||||
.await;
|
||||
return None;
|
||||
}
|
||||
Err(AntigravityV1InternalRequestError::HeaderRulesApplyFailed) => {
|
||||
mark_skipped_local_standard_candidate_with_failure_diagnostic(
|
||||
state,
|
||||
input,
|
||||
trace_id,
|
||||
candidate,
|
||||
attempt.candidate_index,
|
||||
&attempt.candidate_id,
|
||||
"transport_header_rules_apply_failed",
|
||||
CandidateFailureDiagnostic::header_rules_apply_failed(
|
||||
client_api_format,
|
||||
provider_api_format,
|
||||
"standard_family_antigravity_headers",
|
||||
),
|
||||
)
|
||||
.await;
|
||||
return None;
|
||||
}
|
||||
};
|
||||
|
||||
let mut provider_request_headers = resolved.headers.headers;
|
||||
apply_codex_openai_special_headers(
|
||||
&mut provider_request_headers,
|
||||
&resolved.body,
|
||||
effective_headers,
|
||||
resolved.transport.provider.provider_type.as_str(),
|
||||
provider_api_format,
|
||||
Some(trace_id),
|
||||
resolved.transport.key.decrypted_auth_config.as_deref(),
|
||||
);
|
||||
request_identity_response_encoding_when_redacted(
|
||||
&mut provider_request_headers,
|
||||
request_redacted,
|
||||
);
|
||||
|
||||
Some(LocalStandardCandidatePayloadParts {
|
||||
auth_header: resolved.headers.auth_header,
|
||||
auth_value: resolved.headers.auth_value,
|
||||
mapped_model,
|
||||
provider_api_format: provider_api_format.to_string(),
|
||||
provider_request_body: resolved.body,
|
||||
provider_request_headers,
|
||||
upstream_url: resolved.upstream_url,
|
||||
upstream_is_stream,
|
||||
envelope_name: Some(ANTIGRAVITY_V1INTERNAL_ENVELOPE_NAME),
|
||||
transport: resolved.transport,
|
||||
transport_profile: None,
|
||||
request_redacted,
|
||||
})
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn build_gemini_cli_cross_format_payload_parts(
|
||||
state: &AppState,
|
||||
|
||||
@@ -63,7 +63,7 @@ use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||
};
|
||||
use aether_model_fetch::{
|
||||
aggregate_models_for_cache, fetch_models_from_transports, json_string_list,
|
||||
aggregate_models_for_cache, fetch_models_from_transports_for_management, json_string_list,
|
||||
model_catalog_upstream_metadata, preset_models_for_provider, selected_models_fetch_endpoints,
|
||||
upstream_metadata_namespace_updates,
|
||||
};
|
||||
@@ -505,10 +505,34 @@ async fn provider_query_fetch_models_for_key(
|
||||
key: &StoredProviderCatalogKey,
|
||||
force_refresh: bool,
|
||||
) -> Result<ProviderQueryKeyFetchResult, GatewayError> {
|
||||
let is_codex = provider.provider_type.trim().eq_ignore_ascii_case("codex");
|
||||
let codex_catalog = if is_codex {
|
||||
crate::model_fetch::read_codex_management_catalog(state.app(), &provider.id, &key.id).await
|
||||
} else {
|
||||
None
|
||||
};
|
||||
if !force_refresh {
|
||||
if let Some(cached_models) =
|
||||
// Read through the live, credential-scoped directory before the admin cache.
|
||||
// Never reuse the old versionless cache for Codex (including cached presets).
|
||||
let cached_models = if is_codex {
|
||||
codex_catalog
|
||||
.as_ref()
|
||||
.and_then(|catalog| catalog.models.as_ref())
|
||||
.filter(|_| {
|
||||
selected_models_fetch_endpoints(endpoints, key)
|
||||
.iter()
|
||||
.any(|endpoint| endpoint.api_format == "openai:responses")
|
||||
})
|
||||
.map(|models| {
|
||||
aether_model_fetch::project_codex_models_for_legacy_cache([(
|
||||
"openai:responses",
|
||||
models.as_slice(),
|
||||
)])
|
||||
})
|
||||
} else {
|
||||
provider_query_read_cached_models(state, &provider.id, &key.id).await
|
||||
{
|
||||
};
|
||||
if let Some(cached_models) = cached_models {
|
||||
let models = provider_query_filter_models_for_key(provider, key, cached_models);
|
||||
return Ok(ProviderQueryKeyFetchResult {
|
||||
models,
|
||||
@@ -570,30 +594,51 @@ async fn provider_query_fetch_models_for_key(
|
||||
});
|
||||
}
|
||||
|
||||
let outcome = match fetch_models_from_transports(state.app(), &transports).await {
|
||||
Ok(outcome) => outcome,
|
||||
Err(err) => {
|
||||
all_errors.extend(provider_query_project_model_fetch_errors([err]));
|
||||
if let Some(fallback) =
|
||||
provider_query_codex_preset_fallback(provider, &all_errors.join("; "))
|
||||
{
|
||||
provider_query_persist_preset_models(state, provider, key, &fallback.models)
|
||||
.await?;
|
||||
return Ok(fallback);
|
||||
let client_version = is_codex.then(|| {
|
||||
codex_catalog
|
||||
.as_ref()
|
||||
.map(|catalog| catalog.client_version.as_str())
|
||||
.unwrap_or(crate::ai_serving::CODEX_CLIENT_VERSION)
|
||||
});
|
||||
let outcome =
|
||||
match fetch_models_from_transports_for_management(state.app(), &transports, client_version)
|
||||
.await
|
||||
{
|
||||
Ok(outcome) => outcome,
|
||||
Err(err) => {
|
||||
all_errors.extend(provider_query_project_model_fetch_errors([err]));
|
||||
if let Some(fallback) =
|
||||
provider_query_codex_preset_fallback(provider, &all_errors.join("; "))
|
||||
{
|
||||
provider_query_persist_preset_models(state, provider, key, &fallback.models)
|
||||
.await?;
|
||||
return Ok(fallback);
|
||||
}
|
||||
return Ok(ProviderQueryKeyFetchResult {
|
||||
models: Vec::new(),
|
||||
error: Some(all_errors.join("; ")),
|
||||
warning: None,
|
||||
from_cache: false,
|
||||
has_success: false,
|
||||
});
|
||||
}
|
||||
return Ok(ProviderQueryKeyFetchResult {
|
||||
models: Vec::new(),
|
||||
error: Some(all_errors.join("; ")),
|
||||
warning: None,
|
||||
from_cache: false,
|
||||
has_success: false,
|
||||
});
|
||||
}
|
||||
};
|
||||
};
|
||||
|
||||
all_errors.extend(provider_query_project_model_fetch_errors(outcome.errors));
|
||||
let unique_models = outcome.legacy_models;
|
||||
if outcome.has_success && !unique_models.is_empty() {
|
||||
if all_errors.is_empty() && outcome.native_codex_catalog {
|
||||
if let Some(catalog) = codex_catalog.as_ref() {
|
||||
crate::model_fetch::store_codex_management_catalog(
|
||||
state.app(),
|
||||
catalog,
|
||||
&transports,
|
||||
outcome.cached_models,
|
||||
outcome.etag.as_deref(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
<AppState as ModelFetchRuntimeState>::write_upstream_models_cache(
|
||||
state.app(),
|
||||
&provider.id,
|
||||
|
||||
@@ -1494,6 +1494,133 @@ pub(crate) async fn read_recent_codex_catalog_client_version(
|
||||
(!normalized.used_fallback()).then_some(normalized.value)
|
||||
}
|
||||
|
||||
/// Management is not tied to a downstream client's compatibility version. Keep its
|
||||
/// directory at least as new as the built-in fingerprint and successful catalogs.
|
||||
pub(crate) struct CodexManagementCatalog {
|
||||
pub(crate) client_version: String,
|
||||
pub(crate) models: Option<Vec<Value>>,
|
||||
target: CodexCatalogTarget,
|
||||
}
|
||||
|
||||
pub(crate) async fn read_codex_management_catalog<R>(
|
||||
runtime: &R,
|
||||
provider_id: &str,
|
||||
key_id: &str,
|
||||
) -> Option<CodexManagementCatalog>
|
||||
where
|
||||
R: CodexCatalogRuntime + ?Sized,
|
||||
{
|
||||
let target = bind_codex_catalog_target(
|
||||
runtime,
|
||||
&CodexCatalogTarget {
|
||||
identity: CodexCatalogIdentity::new(provider_id, key_id),
|
||||
endpoint_ids: Vec::new(),
|
||||
credential_scope: None,
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
let scope = target.credential_scope()?;
|
||||
let state = runtime.codex_catalog_runtime_state();
|
||||
let mut version = Version::parse(crate::ai_serving::CODEX_CLIENT_VERSION).ok()?;
|
||||
if let Some(recent) =
|
||||
read_recent_codex_catalog_client_version(state, provider_id, key_id, scope).await
|
||||
{
|
||||
version = version.max(Version::parse(&recent).ok()?);
|
||||
}
|
||||
// The most recently seen client can be older than an already successful catalog.
|
||||
// Only consider the current credential generation's bounded success index.
|
||||
for member in state
|
||||
.score_range_by_min(&catalog_versions_key(&target.identity), 0.0)
|
||||
.await
|
||||
.unwrap_or_default()
|
||||
{
|
||||
if let Some((stored_scope, stored_version)) = parse_catalog_version_member(&member) {
|
||||
if stored_scope == scope {
|
||||
if let Ok(candidate) = Version::parse(stored_version) {
|
||||
version = version.max(candidate);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
let client_version = version.to_string();
|
||||
let models = if let Some(snapshot) = read_lkg_snapshot(runtime, &target, &client_version).await
|
||||
{
|
||||
let fresh = state
|
||||
.kv_get(&catalog_fresh_key(&target, &client_version))
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
(fresh.as_deref() == Some(snapshot.content_sha256.as_str())).then_some(snapshot.models)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
if !codex_catalog_credential_scope_is_current(runtime, &target).await {
|
||||
return None;
|
||||
}
|
||||
Some(CodexManagementCatalog {
|
||||
client_version,
|
||||
models,
|
||||
target,
|
||||
})
|
||||
}
|
||||
|
||||
/// Publish a successful management fetch into the same versioned directory used by
|
||||
/// clients. A credential replacement while the request is in flight must not leak
|
||||
/// the previous account's catalog into the new account.
|
||||
pub(crate) async fn store_codex_management_catalog<R>(
|
||||
runtime: &R,
|
||||
catalog: &CodexManagementCatalog,
|
||||
transports: &[GatewayProviderTransportSnapshot],
|
||||
models: Vec<Value>,
|
||||
etag: Option<&str>,
|
||||
) where
|
||||
R: CodexCatalogRuntime + ?Sized,
|
||||
{
|
||||
let target = &catalog.target;
|
||||
if transports.is_empty()
|
||||
|| transports.iter().any(|transport| {
|
||||
codex_catalog_credential_scope_from_transport(transport).as_deref()
|
||||
!= target.credential_scope()
|
||||
})
|
||||
|| !codex_catalog_credential_scope_is_current(runtime, target).await
|
||||
{
|
||||
return;
|
||||
}
|
||||
let Ok(serialized) = validate_catalog_models(&models) else {
|
||||
return;
|
||||
};
|
||||
let now = current_unix_secs();
|
||||
let snapshot = CodexCatalogSnapshot {
|
||||
schema_version: CODEX_CATALOG_SCHEMA_VERSION,
|
||||
provider_id: target.identity.provider_id.clone(),
|
||||
key_id: target.identity.key_id.clone(),
|
||||
credential_scope: target.credential_scope().unwrap_or_default().to_string(),
|
||||
client_version: catalog.client_version.clone(),
|
||||
models,
|
||||
etag: etag.and_then(normalize_etag),
|
||||
fetched_at_unix_secs: now,
|
||||
last_checked_at_unix_secs: now,
|
||||
content_sha256: sha256_hex(&serialized),
|
||||
};
|
||||
persist_catalog_success(
|
||||
runtime,
|
||||
target,
|
||||
&catalog.client_version,
|
||||
&snapshot,
|
||||
Some(200),
|
||||
Duration::ZERO,
|
||||
)
|
||||
.await;
|
||||
if !codex_catalog_credential_scope_is_current(runtime, target).await {
|
||||
discard_catalog_success_after_scope_change(
|
||||
runtime.codex_catalog_runtime_state(),
|
||||
target,
|
||||
&catalog.client_version,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn retain_catalog_version(
|
||||
state: &RuntimeState,
|
||||
target: &CodexCatalogTarget,
|
||||
@@ -2327,6 +2454,101 @@ mod tests {
|
||||
assert!(!prerelease.used_fallback());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn management_catalog_uses_version_floor_and_newest_success_not_last_client() {
|
||||
let runtime = TestRuntime::new(vec![successful_execution("gpt-new", "etag")]);
|
||||
remember_seen_version(&runtime.state, &target(), "0.144.1").await;
|
||||
let initial = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.expect("management context");
|
||||
assert_eq!(
|
||||
initial.client_version,
|
||||
crate::ai_serving::CODEX_CLIENT_VERSION
|
||||
);
|
||||
assert!(initial.models.is_none());
|
||||
|
||||
seed_catalog(&runtime, &version("0.200.0")).await;
|
||||
remember_seen_version(&runtime.state, &target(), "0.144.1").await;
|
||||
let current = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.expect("management context");
|
||||
assert_eq!(current.client_version, "0.200.0");
|
||||
assert_eq!(current.models.unwrap()[0]["slug"], "gpt-new");
|
||||
assert_eq!(
|
||||
runtime.execution_count(),
|
||||
1,
|
||||
"management read does not fetch upstream"
|
||||
);
|
||||
|
||||
runtime
|
||||
.state
|
||||
.kv_delete(&catalog_fresh_key(&target(), "0.200.0"))
|
||||
.await
|
||||
.unwrap();
|
||||
let stale = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stale.client_version, "0.200.0");
|
||||
assert!(
|
||||
stale.models.is_none(),
|
||||
"expired LKG must not become a fresh admin cache"
|
||||
);
|
||||
|
||||
runtime.set_credential_generation(TEST_CREDENTIAL_GENERATION_B);
|
||||
let rebound = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
rebound.client_version,
|
||||
crate::ai_serving::CODEX_CLIENT_VERSION
|
||||
);
|
||||
assert!(rebound.models.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn management_refresh_updates_shared_catalog_but_rejects_replaced_credentials() {
|
||||
let runtime = TestRuntime::new(vec![]);
|
||||
let context = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.unwrap();
|
||||
let transports = vec![sample_codex_transport()];
|
||||
for slug in ["gpt-old", "gpt-6-astra"] {
|
||||
store_codex_management_catalog(
|
||||
&runtime,
|
||||
&context,
|
||||
&transports,
|
||||
vec![codex_model(slug)],
|
||||
Some("etag"),
|
||||
)
|
||||
.await;
|
||||
let shared = read_lkg_snapshot(&runtime, &target(), &context.client_version)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(shared.models[0]["slug"], slug);
|
||||
let admin = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(admin.models.unwrap()[0]["slug"], slug);
|
||||
}
|
||||
runtime.set_credential_generation(TEST_CREDENTIAL_GENERATION_B);
|
||||
store_codex_management_catalog(
|
||||
&runtime,
|
||||
&context,
|
||||
&transports,
|
||||
vec![codex_model("gpt-leaked")],
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
let rebound = read_codex_management_catalog(&runtime, TEST_PROVIDER_ID, TEST_KEY_ID)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(rebound.models.is_none());
|
||||
let old = read_raw_snapshot(&runtime, &target(), &context.client_version)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(old.models[0]["slug"], "gpt-6-astra");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invalid_or_oversized_client_versions_use_bounded_fallback_identity() {
|
||||
for raw in [
|
||||
|
||||
@@ -6,8 +6,8 @@ mod tests;
|
||||
pub(crate) use aether_model_fetch::ModelFetchRunSummary;
|
||||
pub(crate) use catalog::{
|
||||
codex_catalog_credential_scope_from_stored_key, codex_catalog_targets, load_codex_catalogs,
|
||||
normalize_codex_client_version, read_recent_codex_catalog_client_version,
|
||||
refresh_codex_catalog_target, CodexCatalogLoad, CodexCatalogRuntime, CodexCatalogTarget,
|
||||
normalize_codex_client_version, read_codex_management_catalog, refresh_codex_catalog_target,
|
||||
store_codex_management_catalog, CodexCatalogLoad, CodexCatalogRuntime, CodexCatalogTarget,
|
||||
NormalizedCodexClientVersion,
|
||||
};
|
||||
pub(crate) use runtime::state::ModelFetchRuntimeState;
|
||||
|
||||
@@ -7,7 +7,7 @@ use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||
};
|
||||
use aether_model_fetch::{
|
||||
apply_model_filters, fetch_models_from_transports_for_client_version, json_string_list,
|
||||
apply_model_filters, fetch_models_from_transports_for_management, json_string_list,
|
||||
model_catalog_upstream_metadata, model_fetch_interval_minutes,
|
||||
model_fetch_startup_delay_seconds, model_fetch_startup_enabled, preset_models_for_provider,
|
||||
selected_models_fetch_endpoints, sync_provider_model_whitelist_associations,
|
||||
@@ -433,7 +433,7 @@ async fn fetch_and_persist_key_models(
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let result = match fetch_models_from_transports_for_client_version(
|
||||
let result = match fetch_models_from_transports_for_management(
|
||||
state,
|
||||
&transports,
|
||||
codex_client_version.as_deref(),
|
||||
|
||||
@@ -455,22 +455,9 @@ impl ModelFetchRuntimeState for AppState {
|
||||
provider_id: &str,
|
||||
key_id: &str,
|
||||
) -> Option<String> {
|
||||
let credential_scope =
|
||||
<AppState as CodexCatalogRuntime>::read_codex_catalog_credential_scope_strong(
|
||||
self,
|
||||
provider_id,
|
||||
key_id,
|
||||
)
|
||||
crate::model_fetch::read_codex_management_catalog(self, provider_id, key_id)
|
||||
.await
|
||||
.ok()
|
||||
.flatten()?;
|
||||
crate::model_fetch::read_recent_codex_catalog_client_version(
|
||||
self.runtime_state.as_ref(),
|
||||
provider_id,
|
||||
key_id,
|
||||
&credential_scope,
|
||||
)
|
||||
.await
|
||||
.map(|catalog| catalog.client_version)
|
||||
}
|
||||
|
||||
async fn update_provider_catalog_key_model_fetch_state(
|
||||
|
||||
@@ -421,7 +421,7 @@ async fn gateway_executes_codex_image_stream_via_local_decision_gate_after_oauth
|
||||
);
|
||||
assert_eq!(
|
||||
seen_execution_runtime_request.headers["user-agent"],
|
||||
"codex_cli_rs/0.144.1"
|
||||
"codex_cli_rs/0.153.3"
|
||||
);
|
||||
assert_eq!(
|
||||
seen_execution_runtime_request.headers["originator"],
|
||||
|
||||
@@ -1119,7 +1119,7 @@ async fn gateway_executes_codex_image_sync_via_local_decision_gate_after_oauth_r
|
||||
);
|
||||
assert_eq!(
|
||||
seen_execution_runtime_request.headers["user-agent"],
|
||||
"codex_cli_rs/0.144.1"
|
||||
"codex_cli_rs/0.153.3"
|
||||
);
|
||||
assert_eq!(
|
||||
seen_execution_runtime_request.headers["originator"],
|
||||
|
||||
@@ -2896,6 +2896,22 @@ fn ai_serving_standard_attempts_consume_eligible_local_candidates_without_transp
|
||||
);
|
||||
}
|
||||
|
||||
let standard_family_request = read_workspace_file(
|
||||
"apps/aether-gateway/src/ai_serving/planner/standard/family/request.rs",
|
||||
);
|
||||
for pattern in [
|
||||
"is_antigravity_provider_transport(",
|
||||
"build_antigravity_v1internal_provider_request(",
|
||||
"is_gemini_cli_provider_transport(",
|
||||
"build_gemini_cli_v1internal_provider_request(",
|
||||
] {
|
||||
assert!(
|
||||
standard_family_request.contains(pattern),
|
||||
"standard family request preparation should build v1internal envelopes through {pattern} \
|
||||
so a cross-format client never posts a bare Gemini body to a v1internal URL"
|
||||
);
|
||||
}
|
||||
|
||||
let provider_transport_standard =
|
||||
read_workspace_file("crates/aether-provider/transport/src/standard/mod.rs");
|
||||
for pattern in [
|
||||
|
||||
@@ -538,14 +538,14 @@ async fn gateway_handles_admin_provider_query_models_with_openai_responses_endpo
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache() {
|
||||
fn gateway_recovers_codex_slug_only_models_from_a_stale_legacy_cache() {
|
||||
run_provider_query_test(
|
||||
"gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache",
|
||||
gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache_impl,
|
||||
"gateway_recovers_codex_slug_only_models_from_a_stale_legacy_cache",
|
||||
gateway_recovers_codex_slug_only_models_from_a_stale_legacy_cache_impl,
|
||||
);
|
||||
}
|
||||
|
||||
async fn gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache_impl() {
|
||||
async fn gateway_recovers_codex_slug_only_models_from_a_stale_legacy_cache_impl() {
|
||||
let execution_runtime_hits = Arc::new(Mutex::new(0usize));
|
||||
let execution_runtime_hits_clone = Arc::clone(&execution_runtime_hits);
|
||||
let execution_runtime = Router::new().route(
|
||||
@@ -558,7 +558,11 @@ async fn gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache_impl
|
||||
.expect("mutex should lock") += 1;
|
||||
assert_eq!(
|
||||
plan.url,
|
||||
"https://chatgpt.com/backend-api/codex/models?client_version=0.144.1"
|
||||
"https://chatgpt.com/backend-api/codex/models?client_version=0.153.3"
|
||||
);
|
||||
assert_eq!(
|
||||
plan.headers.get("user-agent").map(String::as_str),
|
||||
Some(aether_ai_formats::CODEX_CLIENT_USER_AGENT)
|
||||
);
|
||||
assert_eq!(plan.provider_api_format, "openai:responses");
|
||||
Json(json!({
|
||||
@@ -617,16 +621,21 @@ async fn gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache_impl
|
||||
.runtime_state()
|
||||
.kv_set(
|
||||
"upstream_models:provider-codex-dynamic:key-codex-dynamic",
|
||||
"[]".to_string(),
|
||||
json!([{"id": "gpt-stale-legacy"}]).to_string(),
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.expect("empty legacy cache should seed");
|
||||
.expect("stale legacy cache should seed");
|
||||
let cache_state = state.clone();
|
||||
let gateway = build_router_with_state(state);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
for (request_index, expected_from_cache) in [(0usize, false), (1usize, true)] {
|
||||
for (request_index, expected_from_cache) in [
|
||||
(0usize, false),
|
||||
(1usize, true),
|
||||
(2usize, false),
|
||||
(3usize, true),
|
||||
] {
|
||||
let response = reqwest::Client::new()
|
||||
.post(format!("{gateway_url}/api/admin/provider-query/models"))
|
||||
.header(crate::constants::GATEWAY_HEADER, "rust-phase3b")
|
||||
@@ -635,7 +644,8 @@ async fn gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache_impl
|
||||
.header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123")
|
||||
.json(&json!({
|
||||
"provider_id": "provider-codex-dynamic",
|
||||
"api_key_id": "key-codex-dynamic"
|
||||
"api_key_id": "key-codex-dynamic",
|
||||
"force_refresh": request_index == 2,
|
||||
}))
|
||||
.send()
|
||||
.await
|
||||
@@ -672,10 +682,47 @@ async fn gateway_recovers_codex_slug_only_models_from_an_empty_legacy_cache_impl
|
||||
);
|
||||
assert_eq!(
|
||||
*execution_runtime_hits.lock().expect("mutex should lock"),
|
||||
1,
|
||||
if request_index < 2 { 1 } else { 2 },
|
||||
"request {request_index} must not cause another upstream fetch"
|
||||
);
|
||||
assert_eq!(
|
||||
payload["data"]["models"][0]["display_name"],
|
||||
if request_index == 1 {
|
||||
"Updated shared catalog"
|
||||
} else {
|
||||
"Future Dynamic"
|
||||
}
|
||||
);
|
||||
if request_index == 0 {
|
||||
let context = crate::model_fetch::read_codex_management_catalog(
|
||||
&cache_state,
|
||||
"provider-codex-dynamic",
|
||||
"key-codex-dynamic",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let mut models = context
|
||||
.models
|
||||
.clone()
|
||||
.expect("admin fetch publishes shared catalog");
|
||||
models[0]["display_name"] = json!("Updated shared catalog");
|
||||
let transport = cache_state
|
||||
.read_provider_transport_snapshot(
|
||||
"provider-codex-dynamic",
|
||||
"endpoint-codex-dynamic",
|
||||
"key-codex-dynamic",
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
crate::model_fetch::store_codex_management_catalog(
|
||||
&cache_state,
|
||||
&context,
|
||||
&[transport],
|
||||
models,
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
<AppState as crate::model_fetch::ModelFetchRuntimeState>::write_upstream_models_cache(
|
||||
&cache_state,
|
||||
"provider-codex-dynamic",
|
||||
@@ -712,7 +759,7 @@ async fn gateway_handles_admin_provider_query_models_falls_back_to_codex_preset_
|
||||
.expect("mutex should lock") += 1;
|
||||
assert_eq!(
|
||||
plan.url,
|
||||
"https://chatgpt.com/backend-api/codex/models?client_version=0.144.1"
|
||||
"https://chatgpt.com/backend-api/codex/models?client_version=0.153.3"
|
||||
);
|
||||
Json(json!({
|
||||
"request_id": "req-provider-query-codex-invalidated",
|
||||
|
||||
Reference in New Issue
Block a user