Merge branch 'pr-575'

# Conflicts:
#	apps/aether-gateway/src/ai_serving/planner/passthrough/provider/family/payload.rs
#	apps/aether-gateway/src/ai_serving/planner/passthrough/provider/family/request.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/family/payload.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/family/request.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/chat/decision/request.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/responses/decision/payload.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/responses/decision/request.rs
This commit is contained in:
fawney19
2026-05-27 01:34:19 +08:00
13 changed files with 2052 additions and 348 deletions
@@ -17,6 +17,7 @@ mod passthrough;
mod plan_builders;
mod pool_scheduler;
pub(crate) mod pool_scores;
mod redaction;
mod report_context;
mod route;
mod runtime_miss;
@@ -56,10 +56,15 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
let Some(resolved) = resolve_local_same_format_provider_candidate_payload_parts(
state, parts, trace_id, body_json, input, &attempt, spec,
)
.await
.await?
else {
return Ok(None);
};
let original_request_body_json = if resolved.request_redacted {
Some(&resolved.provider_request_body)
} else {
Some(body_json)
};
let prompt_cache_key = resolved
.provider_request_body
@@ -130,7 +135,7 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
request_path: Some(parts.uri.path()),
request_query_string: parts.uri.query(),
request_origin: Some(crate::ai_serving::request_origin_from_parts(parts)),
original_request_body_json: Some(body_json),
original_request_body_json,
original_request_body_base64: None,
client_session_affinity: input.client_session_affinity.as_ref(),
scheduler_affinity_epoch: eligible.orchestration.scheduler_affinity_epoch,
@@ -164,6 +169,7 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
provider_request_headers,
provider_request_body,
transport_profile: _,
request_redacted: _,
} = resolved;
let mut decision = build_ai_execution_decision_response(AiExecutionDecisionResponseParts {
@@ -7,6 +7,9 @@ use serde_json::Value;
use crate::ai_serving::planner::common::{
enforce_provider_body_stream_policy, request_requires_body_stream_field,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
use crate::ai_serving::transport::antigravity::{
build_antigravity_safe_v1internal_request, build_antigravity_static_identity_headers,
classify_local_antigravity_request_support, AntigravityEnvelopeRequestType,
@@ -17,7 +20,7 @@ use crate::ai_serving::transport::{
GrokHeaderInput, SameFormatProviderHeadersInput, GROK_CHAT_PATH,
};
use crate::ai_serving::{CandidateFailureDiagnostic, GatewayProviderTransportSnapshot};
use crate::AppState;
use crate::{AppState, GatewayError};
mod policy;
mod prepare;
@@ -99,6 +102,7 @@ pub(crate) struct LocalSameFormatProviderCandidatePayloadParts {
pub(super) provider_request_headers: BTreeMap<String, String>,
pub(super) provider_request_body: Value,
pub(super) transport_profile: Option<ResolvedTransportProfile>,
pub(super) request_redacted: bool,
}
pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
@@ -109,9 +113,9 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
input: &LocalSameFormatProviderDecisionInput,
attempt: &LocalSameFormatProviderCandidateAttempt,
spec: LocalSameFormatProviderSpec,
) -> Option<LocalSameFormatProviderCandidatePayloadParts> {
) -> Result<Option<LocalSameFormatProviderCandidatePayloadParts>, GatewayError> {
let candidate = &attempt.eligible.candidate;
let prepared = prepare_local_same_format_provider_candidate(
let Some(prepared) = prepare_local_same_format_provider_candidate(
state,
trace_id,
input,
@@ -120,7 +124,10 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
&attempt.candidate_id,
spec,
)
.await?;
.await
else {
return Ok(None);
};
let enable_model_directives =
crate::system_features::reasoning_model_directive_enabled_for_api_format_and_model(
state,
@@ -129,6 +136,16 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
)
.await;
let effective_headers = input.effective_headers(&parts.headers);
let redaction = resolve_provider_chat_pii_redaction(
state,
parts,
body_json,
&input.auth_context,
spec.api_format,
&attempt.candidate_id,
)
.await?;
let body_json = redaction.body_json.as_ref();
let Some(mut base_provider_request_body) =
super::super::request::build_same_format_provider_request_body(
@@ -165,7 +182,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
if let Some(mapping) =
crate::system_features::reasoning_model_directive_mapping_for_api_format_and_model(
@@ -211,7 +228,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
"transport_unsupported",
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -243,7 +260,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -288,14 +305,14 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
let extra_headers = antigravity_auth
.as_ref()
.map(build_antigravity_static_identity_headers)
.unwrap_or_default();
let Some(provider_request_headers) = (if is_grok {
let Some(mut provider_request_headers) = (if is_grok {
build_grok_browser_headers(GrokHeaderInput {
transport: &prepared.transport,
transport_profile: transport_profile.as_ref(),
@@ -339,10 +356,14 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
redaction.redacted,
);
Some(LocalSameFormatProviderCandidatePayloadParts {
Ok(Some(LocalSameFormatProviderCandidatePayloadParts {
transport: prepared.transport,
is_antigravity: prepared.is_antigravity,
is_kiro: prepared.is_kiro,
@@ -356,5 +377,6 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
provider_request_headers,
provider_request_body,
transport_profile,
})
request_redacted: redaction.redacted,
}))
}
@@ -0,0 +1,192 @@
use std::borrow::Cow;
use std::time::{SystemTime, UNIX_EPOCH};
use serde_json::Value;
use tracing::warn;
use crate::ai_serving::ExecutionRuntimeAuthContext;
use crate::privacy::{
build_redaction_session_config, read_chat_pii_redaction_runtime_config,
try_mask_chat_pii_request_json_with_cache_options, ChatPiiRedactionRequestFormat,
MaskChatRequestOptions, RedactionMaskError, RedactionSessionSlot, RedisRedactionMappingCache,
};
use crate::{AppState, GatewayError};
pub(crate) struct ProviderRequestRedaction<'a> {
pub(crate) body_json: Cow<'a, Value>,
pub(crate) redacted: bool,
}
impl<'a> ProviderRequestRedaction<'a> {
fn disabled(body_json: &'a Value) -> Self {
Self {
body_json: Cow::Borrowed(body_json),
redacted: false,
}
}
}
#[derive(Clone, Copy, Debug, Default)]
struct ChatPiiRedactionFeatureSettings {
enabled: Option<bool>,
inject_model_instruction: Option<bool>,
}
impl ChatPiiRedactionFeatureSettings {
fn merge_from_value(&mut self, value: Option<&Value>) {
let Some(settings) = value
.and_then(Value::as_object)
.and_then(|features| features.get("chat_pii_redaction"))
.and_then(Value::as_object)
else {
return;
};
if let Some(enabled) = settings.get("enabled").and_then(Value::as_bool) {
self.enabled = Some(enabled);
}
if let Some(inject_model_instruction) = settings
.get("inject_model_instruction")
.and_then(Value::as_bool)
{
self.inject_model_instruction = Some(inject_model_instruction);
}
}
fn effective_enabled(self) -> bool {
self.enabled.unwrap_or(false)
}
fn effective_inject_model_instruction(self) -> bool {
self.inject_model_instruction.unwrap_or(true)
}
}
pub(crate) fn request_identity_response_encoding_when_redacted(
headers: &mut std::collections::BTreeMap<String, String>,
redacted: bool,
) {
if redacted {
headers.insert("accept-encoding".to_string(), "identity".to_string());
}
}
pub(crate) async fn resolve_provider_chat_pii_redaction<'a>(
state: &AppState,
parts: &http::request::Parts,
body_json: &'a Value,
auth_context: &ExecutionRuntimeAuthContext,
client_api_format: &str,
candidate_id: &str,
) -> Result<ProviderRequestRedaction<'a>, GatewayError> {
let Some(format) = ChatPiiRedactionRequestFormat::from_api_format(client_api_format) else {
return Ok(ProviderRequestRedaction::disabled(body_json));
};
let Some(slot) = parts.extensions.get::<RedactionSessionSlot>() else {
return Ok(ProviderRequestRedaction::disabled(body_json));
};
let runtime_config = read_chat_pii_redaction_runtime_config(state)
.await
.map_err(|err| {
warn!(
error = ?err,
"gateway failed to read chat pii redaction runtime config"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
if !runtime_config.enabled {
return Ok(ProviderRequestRedaction::disabled(body_json));
}
let feature_settings = resolve_chat_pii_redaction_feature_settings(state, auth_context).await?;
if !feature_settings.effective_enabled() {
return Ok(ProviderRequestRedaction::disabled(body_json));
}
let Some(hmac_key) = state.encryption_key().map(str::as_bytes).map(Vec::from) else {
warn!("gateway chat pii redaction is enabled but encryption key is unavailable");
return Err(GatewayError::Internal(
"chat pii redaction setup failed".to_string(),
));
};
let body_bytes = serde_json::to_vec(body_json).map_err(|err| {
warn!(
error = ?err,
"gateway failed to serialize provider chat pii redaction body"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
let now_unix_secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let cache = RedisRedactionMappingCache::new(state.runtime_state.as_ref());
let masked = try_mask_chat_pii_request_json_with_cache_options(
&body_bytes,
format,
build_redaction_session_config(hmac_key, &runtime_config, now_unix_secs),
MaskChatRequestOptions::runtime(feature_settings.effective_inject_model_instruction()),
Some(&cache),
)
.await
.map_err(redaction_mask_error_to_gateway_error)?;
if !masked.redacted {
return Ok(ProviderRequestRedaction {
body_json: Cow::Borrowed(body_json),
redacted: false,
});
}
let masked_body_json = serde_json::from_slice::<Value>(&masked.body).map_err(|err| {
warn!(
error = ?err,
"gateway failed to decode redacted provider chat pii body"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
slot.put_for_candidate(candidate_id, masked.session);
Ok(ProviderRequestRedaction {
body_json: Cow::Owned(masked_body_json),
redacted: true,
})
}
async fn resolve_chat_pii_redaction_feature_settings(
state: &AppState,
auth_context: &ExecutionRuntimeAuthContext,
) -> Result<ChatPiiRedactionFeatureSettings, GatewayError> {
let user_settings = state
.read_user_feature_settings(&auth_context.user_id)
.await
.map_err(|err| {
warn!(
error = ?err,
"gateway failed to read user chat pii redaction feature settings"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
let key_settings = state
.read_auth_api_key_feature_settings(
&auth_context.user_id,
&auth_context.api_key_id,
auth_context.api_key_is_standalone,
)
.await
.map_err(|err| {
warn!(
error = ?err,
"gateway failed to read api key chat pii redaction feature settings"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
let mut settings = ChatPiiRedactionFeatureSettings::default();
settings.merge_from_value(user_settings.as_ref());
settings.merge_from_value(key_settings.as_ref());
Ok(settings)
}
fn redaction_mask_error_to_gateway_error(error: RedactionMaskError) -> GatewayError {
match error {
RedactionMaskError::Limit(limit) => GatewayError::Client {
status: limit.client_status(),
message: limit.safe_message().to_string(),
},
}
}
@@ -75,10 +75,15 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
let Some(resolved) = resolve_local_standard_candidate_payload_parts(
state, parts, trace_id, body_json, input, &attempt, spec,
)
.await
.await?
else {
return Ok(None);
};
let original_request_body_json = if resolved.request_redacted {
Some(&resolved.provider_request_body)
} else {
Some(body_json)
};
let proxy = state
.resolve_transport_proxy_snapshot_with_tunnel_affinity(&resolved.transport)
.await;
@@ -131,7 +136,7 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
request_path: Some(parts.uri.path()),
request_query_string: parts.uri.query(),
request_origin: Some(crate::ai_serving::request_origin_from_parts(parts)),
original_request_body_json: Some(body_json),
original_request_body_json,
original_request_body_base64: None,
client_session_affinity: input.client_session_affinity.as_ref(),
scheduler_affinity_epoch: eligible.orchestration.scheduler_affinity_epoch,
@@ -168,6 +173,7 @@ pub(super) async fn maybe_build_local_standard_decision_payload_for_candidate(
envelope_name: _,
transport,
transport_profile: _,
request_redacted: _,
} = resolved;
let mut decision = build_ai_execution_decision_response(AiExecutionDecisionResponseParts {
@@ -12,6 +12,9 @@ use crate::ai_serving::planner::common::{
endpoint_config_forces_body_stream_field, enforce_provider_body_stream_policy,
request_requires_body_stream_field, resolve_upstream_is_stream_for_provider,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
use crate::ai_serving::planner::spec_metadata::local_standard_spec_metadata;
use crate::ai_serving::planner::standard::{
apply_codex_openai_responses_special_headers, apply_deepseek_tool_call_thinking_compat,
@@ -37,7 +40,7 @@ use crate::ai_serving::{
build_openai_image_request_body_from_gemini_image_request, gemini_request_is_image_generation,
CandidateFailureDiagnostic, GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth,
};
use crate::AppState;
use crate::{AppState, GatewayError};
use super::payload::{
mark_skipped_local_standard_candidate, mark_skipped_local_standard_candidate_with_extra_data,
@@ -59,6 +62,7 @@ pub(crate) struct LocalStandardCandidatePayloadParts {
pub(super) envelope_name: Option<&'static str>,
pub(super) transport: Arc<GatewayProviderTransportSnapshot>,
pub(super) transport_profile: Option<ResolvedTransportProfile>,
pub(super) request_redacted: bool,
}
fn is_grok_text_provider_api_format(provider_api_format: &str) -> bool {
@@ -284,7 +288,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
input: &LocalStandardDecisionInput,
attempt: &LocalStandardCandidateAttempt,
spec: LocalStandardSpec,
) -> Option<LocalStandardCandidatePayloadParts> {
) -> Result<Option<LocalStandardCandidatePayloadParts>, GatewayError> {
let spec_metadata = local_standard_spec_metadata(spec);
let planner_state = crate::ai_serving::PlannerAppState::new(state);
let candidate = &attempt.eligible.candidate;
@@ -301,10 +305,12 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
&& provider_api_format == "openai:image"
&& gemini_request_is_image_generation(body_json)
{
return resolve_local_gemini_image_to_openai_image_candidate_payload_parts(
state, parts, trace_id, body_json, input, attempt,
)
.await;
return Ok(
resolve_local_gemini_image_to_openai_image_candidate_payload_parts(
state, parts, trace_id, body_json, input, attempt,
)
.await,
);
}
let is_kiro_claude_cli = is_kiro_claude_messages_transport(transport, provider_api_format);
if is_grok && is_grok_text_provider_api_format(provider_api_format) {
@@ -333,10 +339,21 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
};
let redaction = resolve_provider_chat_pii_redaction(
state,
parts,
body_json,
&input.auth_context,
spec_metadata.api_format,
&attempt.candidate_id,
)
.await?;
let body_json = redaction.body_json.as_ref();
let mut provider_request_body = body_json.clone();
if let Some(object) = provider_request_body.as_object_mut() {
object.insert(
@@ -362,7 +379,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
);
let upstream_url = build_grok_upstream_url(transport, GROK_CHAT_PATH);
let Some(provider_request_headers) = build_grok_browser_headers(GrokHeaderInput {
let Some(mut provider_request_headers) = build_grok_browser_headers(GrokHeaderInput {
transport,
transport_profile: transport_profile.as_ref(),
request_headers: Some(effective_headers),
@@ -387,10 +404,14 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
redaction.redacted,
);
return Some(LocalStandardCandidatePayloadParts {
return Ok(Some(LocalStandardCandidatePayloadParts {
auth_header: prepared_candidate.auth_header,
auth_value: prepared_candidate.auth_value,
mapped_model: prepared_candidate.mapped_model,
@@ -402,7 +423,8 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
envelope_name: None,
transport: Arc::clone(transport),
transport_profile,
});
request_redacted: redaction.redacted,
}));
}
if !crate::ai_serving::request_pair_allowed_for_transport(
@@ -410,7 +432,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
spec_metadata.api_format,
provider_api_format,
) {
return None;
return Ok(None);
}
let is_windsurf_cascade =
@@ -435,7 +457,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
let oauth_context = OauthPreparationContext {
@@ -463,7 +485,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
"transport_auth_unavailable",
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -488,7 +510,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -513,7 +535,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
}
};
@@ -534,6 +556,16 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
Some(&input.requested_model),
)
.await;
let redaction = resolve_provider_chat_pii_redaction(
state,
parts,
body_json,
&input.auth_context,
spec_metadata.api_format,
&attempt.candidate_id,
)
.await?;
let body_json = redaction.body_json.as_ref();
let mut provider_request_body =
match crate::ai_serving::planner::standard::build_standard_request_body_with_model_directives_and_request_headers(
body_json,
@@ -569,7 +601,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
};
enforce_provider_body_stream_policy(
@@ -599,7 +631,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
apply_non_native_claude_thinking_signature_compat(
&mut provider_request_body,
@@ -654,7 +686,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
apply_non_native_claude_thinking_signature_compat(
&mut provider_request_body,
@@ -671,7 +703,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
}
if let Some(kiro_auth) = kiro_auth.as_ref() {
return build_kiro_cross_format_payload_parts(
return Ok(build_kiro_cross_format_payload_parts(
state,
parts,
trace_id,
@@ -686,11 +718,12 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
provider_request_body,
upstream_is_stream,
kiro_auth,
redaction.redacted,
)
.await;
.await);
}
if is_windsurf_cascade {
return build_windsurf_cross_format_payload_parts(
return Ok(build_windsurf_cross_format_payload_parts(
state,
parts,
trace_id,
@@ -704,8 +737,9 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
prepared_candidate.auth_value,
provider_request_body,
upstream_is_stream,
redaction.redacted,
)
.await;
.await);
}
let upstream_url = match crate::ai_serving::planner::standard::build_standard_upstream_url(
@@ -733,7 +767,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
};
let Some(resolved_headers) =
@@ -766,7 +800,7 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
let mut provider_request_headers = resolved_headers.headers;
apply_codex_openai_responses_special_headers(
@@ -778,8 +812,12 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
Some(trace_id),
transport.key.decrypted_auth_config.as_deref(),
);
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
redaction.redacted,
);
Some(LocalStandardCandidatePayloadParts {
Ok(Some(LocalStandardCandidatePayloadParts {
auth_header: resolved_headers.auth_header,
auth_value: resolved_headers.auth_value,
mapped_model: prepared_candidate.mapped_model,
@@ -791,7 +829,8 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
envelope_name: None,
transport: Arc::clone(transport),
transport_profile: None,
})
request_redacted: redaction.redacted,
}))
}
fn apply_transport_request_body_semantics(
@@ -821,6 +860,7 @@ async fn build_windsurf_cross_format_payload_parts(
auth_value: String,
openai_chat_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);
@@ -876,7 +916,7 @@ async fn build_windsurf_cross_format_payload_parts(
return None;
}
};
let provider_request_headers = match build_windsurf_cascade_headers(
let mut provider_request_headers = match build_windsurf_cascade_headers(
effective_headers,
&provider_request_body,
original_body_json,
@@ -905,6 +945,10 @@ async fn build_windsurf_cross_format_payload_parts(
return None;
}
};
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
Some(LocalStandardCandidatePayloadParts {
auth_header,
@@ -918,6 +962,7 @@ async fn build_windsurf_cross_format_payload_parts(
envelope_name: Some(WINDSURF_ENVELOPE_NAME),
transport: Arc::clone(transport),
transport_profile: None,
request_redacted,
})
}
@@ -1056,6 +1101,7 @@ async fn resolve_local_gemini_image_to_openai_image_candidate_payload_parts(
envelope_name: None,
transport: Arc::clone(transport),
transport_profile: None,
request_redacted: false,
})
}
@@ -1075,6 +1121,7 @@ async fn build_kiro_cross_format_payload_parts(
claude_request_body: Value,
upstream_is_stream: bool,
kiro_auth: &KiroRequestAuth,
request_redacted: bool,
) -> Option<LocalStandardCandidatePayloadParts> {
let candidate = &attempt.eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
@@ -1133,7 +1180,7 @@ async fn build_kiro_cross_format_payload_parts(
return None;
}
};
let provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
let mut provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
headers: effective_headers,
provider_request_body: &provider_request_body,
original_request_body: original_body_json,
@@ -1163,6 +1210,10 @@ async fn build_kiro_cross_format_payload_parts(
return None;
}
};
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
Some(LocalStandardCandidatePayloadParts {
auth_header,
@@ -1176,6 +1227,7 @@ async fn build_kiro_cross_format_payload_parts(
envelope_name: Some(KIRO_ENVELOPE_NAME),
transport: Arc::clone(transport),
transport_profile: None,
request_redacted,
})
}
@@ -1,7 +1,5 @@
use std::borrow::Cow;
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use aether_contracts::ResolvedTransportProfile;
use serde_json::{json, Value};
@@ -15,6 +13,9 @@ use crate::ai_serving::planner::common::{
endpoint_config_forces_body_stream_field, enforce_provider_body_stream_policy,
request_requires_body_stream_field, OPENAI_CHAT_STREAM_PLAN_KIND,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
use crate::ai_serving::planner::standard::{
apply_codex_openai_responses_special_body_edits, apply_codex_openai_responses_special_headers,
apply_deepseek_tool_call_thinking_compat, build_cross_format_openai_chat_request_body,
@@ -47,13 +48,7 @@ use crate::ai_serving::{
LocalResolvedOAuthRequestAuth,
};
use crate::ai_serving::{ConversionMode, ExecutionStrategy};
use crate::privacy::{
build_redaction_session_config, read_chat_pii_redaction_runtime_config,
try_mask_chat_request_json_with_cache_options, MaskChatRequestOptions, RedactionMaskError,
RedactionSessionSlot, RedisRedactionMappingCache,
};
use crate::{AppState, GatewayError};
use tracing::warn;
use super::support::{
mark_skipped_local_openai_chat_candidate,
@@ -87,100 +82,6 @@ fn is_grok_text_provider_api_format(provider_api_format: &str) -> bool {
)
}
fn request_identity_response_encoding_when_redacted(
headers: &mut BTreeMap<String, String>,
redacted: bool,
) {
if redacted {
headers.insert("accept-encoding".to_string(), "identity".to_string());
}
}
struct ProviderChatRequestRedaction<'a> {
body_json: Cow<'a, Value>,
redacted: bool,
}
impl<'a> ProviderChatRequestRedaction<'a> {
fn disabled(body_json: &'a Value, _parts: &http::request::Parts) -> Self {
Self {
body_json: Cow::Borrowed(body_json),
redacted: false,
}
}
}
#[derive(Clone, Copy, Debug, Default)]
struct ChatPiiRedactionFeatureSettings {
enabled: Option<bool>,
inject_model_instruction: Option<bool>,
}
impl ChatPiiRedactionFeatureSettings {
fn merge_from_value(&mut self, value: Option<&Value>) {
let Some(settings) = value
.and_then(Value::as_object)
.and_then(|features| features.get("chat_pii_redaction"))
.and_then(Value::as_object)
else {
return;
};
if let Some(enabled) = settings.get("enabled").and_then(Value::as_bool) {
self.enabled = Some(enabled);
}
if let Some(inject_model_instruction) = settings
.get("inject_model_instruction")
.and_then(Value::as_bool)
{
self.inject_model_instruction = Some(inject_model_instruction);
}
}
fn effective_enabled(self) -> bool {
self.enabled.unwrap_or(false)
}
fn effective_inject_model_instruction(self) -> bool {
self.inject_model_instruction.unwrap_or(true)
}
}
async fn resolve_chat_pii_redaction_feature_settings(
state: &AppState,
input: &LocalOpenAiChatDecisionInput,
) -> Result<ChatPiiRedactionFeatureSettings, GatewayError> {
let user_settings = state
.read_user_feature_settings(&input.auth_context.user_id)
.await
.map_err(|err| {
warn!(
error = ?err,
"gateway failed to read user chat pii redaction feature settings"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
let key_settings = state
.read_auth_api_key_feature_settings(
&input.auth_context.user_id,
&input.auth_context.api_key_id,
input.auth_context.api_key_is_standalone,
)
.await
.map_err(|err| {
warn!(
error = ?err,
"gateway failed to read api key chat pii redaction feature settings"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
let mut settings = ChatPiiRedactionFeatureSettings::default();
settings.merge_from_value(user_settings.as_ref());
settings.merge_from_value(key_settings.as_ref());
Ok(settings)
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
state: &AppState,
@@ -209,9 +110,15 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
Some(&input.requested_model),
)
.await;
let redaction =
resolve_provider_chat_request_redaction(state, parts, body_json, input, candidate_id)
.await?;
let redaction = resolve_provider_chat_pii_redaction(
state,
parts,
body_json,
&input.auth_context,
"openai:chat",
candidate_id,
)
.await?;
let body_json = redaction.body_json.as_ref();
let effective_headers = input.effective_headers(&parts.headers);
let is_grok = transport
@@ -1607,90 +1514,6 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
})
}
async fn resolve_provider_chat_request_redaction<'a>(
state: &AppState,
parts: &http::request::Parts,
body_json: &'a Value,
input: &LocalOpenAiChatDecisionInput,
candidate_id: &str,
) -> Result<ProviderChatRequestRedaction<'a>, GatewayError> {
if parts.uri.path() != "/v1/chat/completions" {
return Ok(ProviderChatRequestRedaction::disabled(body_json, parts));
}
let Some(slot) = parts.extensions.get::<RedactionSessionSlot>() else {
return Ok(ProviderChatRequestRedaction::disabled(body_json, parts));
};
let runtime_config = read_chat_pii_redaction_runtime_config(state)
.await
.map_err(|err| {
warn!(
error = ?err,
"gateway failed to read chat pii redaction runtime config"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
if !runtime_config.enabled {
return Ok(ProviderChatRequestRedaction::disabled(body_json, parts));
}
let feature_settings = resolve_chat_pii_redaction_feature_settings(state, input).await?;
if !feature_settings.effective_enabled() {
return Ok(ProviderChatRequestRedaction::disabled(body_json, parts));
}
let Some(hmac_key) = state.encryption_key().map(str::as_bytes).map(Vec::from) else {
warn!("gateway chat pii redaction is enabled but encryption key is unavailable");
return Err(GatewayError::Internal(
"chat pii redaction setup failed".to_string(),
));
};
let body_bytes = serde_json::to_vec(body_json).map_err(|err| {
warn!(
error = ?err,
"gateway failed to serialize provider chat pii redaction body"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
let now_unix_secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let cache = RedisRedactionMappingCache::new(state.runtime_state.as_ref());
let masked = try_mask_chat_request_json_with_cache_options(
&body_bytes,
build_redaction_session_config(hmac_key, &runtime_config, now_unix_secs),
MaskChatRequestOptions::runtime(feature_settings.effective_inject_model_instruction()),
Some(&cache),
)
.await
.map_err(redaction_mask_error_to_gateway_error)?;
if !masked.redacted {
return Ok(ProviderChatRequestRedaction {
body_json: Cow::Borrowed(body_json),
redacted: false,
});
}
let masked_body_json = serde_json::from_slice::<Value>(&masked.body).map_err(|err| {
warn!(
error = ?err,
"gateway failed to decode redacted provider chat pii body"
);
GatewayError::Internal("chat pii redaction setup failed".to_string())
})?;
slot.put_for_candidate(candidate_id, masked.session);
Ok(ProviderChatRequestRedaction {
body_json: Cow::Owned(masked_body_json),
redacted: true,
})
}
fn redaction_mask_error_to_gateway_error(error: RedactionMaskError) -> GatewayError {
match error {
RedactionMaskError::Limit(limit) => GatewayError::Client {
status: limit.client_status(),
message: limit.safe_message().to_string(),
},
}
}
#[cfg(test)]
mod tests {
use super::*;
@@ -51,11 +51,16 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
&candidate_id,
spec,
)
.await
.await?
else {
return Ok(None);
};
let candidate = &eligible.candidate;
let original_request_body_json = if resolved.request_redacted {
Some(&resolved.provider_request_body)
} else {
Some(body_json)
};
let prompt_cache_key = resolved
.provider_request_body
@@ -142,7 +147,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
request_path: Some(parts.uri.path()),
request_query_string: parts.uri.query(),
request_origin: Some(crate::ai_serving::request_origin_from_parts(parts)),
original_request_body_json: Some(body_json),
original_request_body_json,
original_request_body_base64: None,
client_session_affinity: input.client_session_affinity.as_ref(),
scheduler_affinity_epoch: eligible.orchestration.scheduler_affinity_epoch,
@@ -205,6 +210,7 @@ pub(crate) async fn maybe_build_local_openai_responses_decision_payload_for_cand
transport,
transport_profile: _,
image_request_summary: _,
request_redacted: _,
} = resolved;
let mut decision = build_ai_execution_decision_response(AiExecutionDecisionResponseParts {
@@ -14,6 +14,9 @@ use crate::ai_serving::planner::common::{
endpoint_config_forces_body_stream_field, enforce_provider_body_stream_policy,
request_requires_body_stream_field, resolve_upstream_is_stream_for_provider,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
use crate::ai_serving::planner::spec_metadata::local_openai_responses_spec_metadata;
use crate::ai_serving::planner::standard::{
apply_codex_openai_responses_special_body_edits, apply_codex_openai_responses_special_headers,
@@ -53,7 +56,7 @@ use crate::ai_serving::{
LocalResolvedOAuthRequestAuth, PlannerAppState,
};
use crate::ai_serving::{ConversionMode, ExecutionStrategy};
use crate::AppState;
use crate::{AppState, GatewayError};
use super::support::{
mark_skipped_local_openai_responses_candidate,
@@ -88,6 +91,7 @@ pub(crate) struct LocalOpenAiResponsesCandidatePayloadParts {
pub(super) transport: Arc<GatewayProviderTransportSnapshot>,
pub(super) transport_profile: Option<ResolvedTransportProfile>,
pub(super) image_request_summary: Option<Value>,
pub(super) request_redacted: bool,
}
#[allow(clippy::too_many_arguments)]
@@ -101,7 +105,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
candidate_index: u32,
candidate_id: &str,
spec: LocalOpenAiResponsesSpec,
) -> Option<LocalOpenAiResponsesCandidatePayloadParts> {
) -> Result<Option<LocalOpenAiResponsesCandidatePayloadParts>, GatewayError> {
let spec_metadata = local_openai_responses_spec_metadata(spec);
let client_api_format = spec_metadata.api_format.trim().to_ascii_lowercase();
let planner_state = PlannerAppState::new(state);
@@ -118,7 +122,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
.eq_ignore_ascii_case("grok");
if !is_grok && provider_api_format.eq_ignore_ascii_case("openai:image") {
return resolve_openai_responses_to_openai_image_payload_parts(
return Ok(resolve_openai_responses_to_openai_image_payload_parts(
state,
parts,
trace_id,
@@ -129,7 +133,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
candidate_id,
spec,
)
.await;
.await);
}
let is_windsurf_cascade =
provider_api_format == "openai:chat" && is_windsurf_provider_transport(transport);
@@ -166,7 +170,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
let oauth_context = OauthPreparationContext {
@@ -194,7 +198,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
"transport_auth_unavailable",
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -235,7 +239,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -260,7 +264,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
skip_reason,
)
.await;
return None;
return Ok(None);
}
}
};
@@ -274,6 +278,16 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
Some(&input.requested_model),
)
.await;
let redaction = resolve_provider_chat_pii_redaction(
state,
parts,
body_json,
&input.auth_context,
spec_metadata.api_format,
candidate_id,
)
.await?;
let body_json = redaction.body_json.as_ref();
let needs_bidirectional_conversion = !same_format && conversion_kind.is_some();
let upstream_is_stream = resolve_upstream_is_stream_for_provider(
@@ -352,7 +366,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
if let Some(mapping) =
crate::system_features::reasoning_model_directive_mapping_for_api_format_and_model(
@@ -400,7 +414,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
"transport_unsupported",
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -431,7 +445,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
}
} else {
@@ -458,11 +472,12 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
upstream_is_stream,
needs_bidirectional_conversion,
kiro_auth,
redaction.redacted,
)
.await;
}
if is_windsurf_cascade {
return build_windsurf_openai_responses_payload_parts(
return Ok(build_windsurf_openai_responses_payload_parts(
state,
parts,
trace_id,
@@ -479,8 +494,9 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
auth_value,
provider_request_body,
upstream_is_stream,
redaction.redacted,
)
.await;
.await);
}
let Some(upstream_url) = (if is_grok && is_grok_text_provider_api_format(provider_api_format) {
@@ -516,7 +532,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
let extra_headers = antigravity_auth
.as_ref()
@@ -548,7 +564,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
crate::ai_serving::transport::StandardProviderRequestHeaders {
headers,
@@ -586,7 +602,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
),
)
.await;
return None;
return Ok(None);
};
resolved_headers
};
@@ -602,6 +618,10 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
transport.key.decrypted_auth_config.as_deref(),
);
}
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
redaction.redacted,
);
let (execution_strategy, conversion_mode) =
ai_local_execution_contract_for_formats(spec_metadata.api_format, provider_api_format);
@@ -630,7 +650,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
"gateway resolved local openai responses upstream url"
);
Some(LocalOpenAiResponsesCandidatePayloadParts {
Ok(Some(LocalOpenAiResponsesCandidatePayloadParts {
auth_header: resolved_headers.auth_header,
auth_value: resolved_headers.auth_value,
mapped_model,
@@ -651,7 +671,8 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
transport: Arc::clone(transport),
transport_profile,
image_request_summary: None,
})
request_redacted: redaction.redacted,
}))
}
#[allow(clippy::too_many_arguments)]
@@ -672,6 +693,7 @@ async fn build_windsurf_openai_responses_payload_parts(
auth_value: String,
openai_chat_request_body: Value,
upstream_is_stream: bool,
request_redacted: bool,
) -> Option<LocalOpenAiResponsesCandidatePayloadParts> {
let candidate = &eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
@@ -727,7 +749,7 @@ async fn build_windsurf_openai_responses_payload_parts(
return None;
}
};
let provider_request_headers = match build_windsurf_cascade_headers(
let mut provider_request_headers = match build_windsurf_cascade_headers(
effective_headers,
&provider_request_body,
original_body_json,
@@ -756,6 +778,10 @@ async fn build_windsurf_openai_responses_payload_parts(
return None;
}
};
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
let (execution_strategy, conversion_mode) =
ai_local_execution_contract_for_formats(client_api_format, provider_api_format);
@@ -775,6 +801,7 @@ async fn build_windsurf_openai_responses_payload_parts(
transport: Arc::clone(transport),
transport_profile: None,
image_request_summary: None,
request_redacted,
})
}
@@ -964,6 +991,7 @@ async fn resolve_openai_responses_to_openai_image_payload_parts(
transport: Arc::clone(transport),
transport_profile: None,
image_request_summary: Some(image_request_summary),
request_redacted: false,
})
}
@@ -1287,7 +1315,8 @@ async fn build_kiro_openai_responses_payload_parts(
upstream_is_stream: bool,
needs_bidirectional_conversion: bool,
kiro_auth: &KiroRequestAuth,
) -> Option<LocalOpenAiResponsesCandidatePayloadParts> {
request_redacted: bool,
) -> Result<Option<LocalOpenAiResponsesCandidatePayloadParts>, GatewayError> {
let candidate = &eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
let provider_request_body = match build_kiro_provider_request_body(
@@ -1314,7 +1343,7 @@ async fn build_kiro_openai_responses_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
};
let upstream_url = match build_kiro_cross_format_upstream_url(
@@ -1342,10 +1371,10 @@ async fn build_kiro_openai_responses_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
};
let provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
let mut provider_request_headers = match build_kiro_provider_headers(KiroProviderHeadersInput {
headers: effective_headers,
provider_request_body: &provider_request_body,
original_request_body: original_body_json,
@@ -1372,7 +1401,7 @@ async fn build_kiro_openai_responses_payload_parts(
),
)
.await;
return None;
return Ok(None);
}
};
let (execution_strategy, conversion_mode) =
@@ -1397,7 +1426,12 @@ async fn build_kiro_openai_responses_payload_parts(
"gateway resolved local openai responses kiro upstream url"
);
Some(LocalOpenAiResponsesCandidatePayloadParts {
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
Ok(Some(LocalOpenAiResponsesCandidatePayloadParts {
auth_header,
auth_value,
mapped_model,
@@ -1413,7 +1447,8 @@ async fn build_kiro_openai_responses_payload_parts(
transport: Arc::clone(transport),
transport_profile: None,
image_request_summary: None,
})
request_redacted,
}))
}
#[cfg(test)]