mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-09 12:40:20 +08:00
feat(codex): stabilize identity across retries
This commit is contained in:
@@ -0,0 +1,249 @@
|
||||
use std::sync::{Arc, OnceLock};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use aether_provider_transport::CodexFingerprintConvergenceContext;
|
||||
use http::{request::Parts, HeaderMap};
|
||||
use serde_json::Value;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::client_session_affinity::codex_request_signals_from_request;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub(crate) struct CodexFingerprintContextSlot(Arc<OnceLock<CodexFingerprintConvergenceContext>>);
|
||||
|
||||
impl Default for CodexFingerprintContextSlot {
|
||||
fn default() -> Self {
|
||||
Self(Arc::new(OnceLock::new()))
|
||||
}
|
||||
}
|
||||
|
||||
impl CodexFingerprintContextSlot {
|
||||
fn resolve(
|
||||
&self,
|
||||
headers: &HeaderMap,
|
||||
body_json: &Value,
|
||||
) -> CodexFingerprintConvergenceContext {
|
||||
self.0
|
||||
.get_or_init(|| {
|
||||
build_codex_fingerprint_context(headers, body_json, Uuid::now_v7().to_string())
|
||||
})
|
||||
.clone()
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn resolve_codex_fingerprint_context(
|
||||
parts: &Parts,
|
||||
body_json: &Value,
|
||||
) -> CodexFingerprintConvergenceContext {
|
||||
if let Some(context) = parts
|
||||
.extensions
|
||||
.get::<CodexFingerprintConvergenceContext>()
|
||||
.cloned()
|
||||
{
|
||||
return context;
|
||||
}
|
||||
if let Some(slot) = parts.extensions.get::<CodexFingerprintContextSlot>() {
|
||||
return slot.resolve(&parts.headers, body_json);
|
||||
}
|
||||
build_codex_fingerprint_context(&parts.headers, body_json, Uuid::now_v7().to_string())
|
||||
}
|
||||
|
||||
pub(crate) fn install_codex_fingerprint_context_slot(parts: &mut Parts) {
|
||||
if parts
|
||||
.extensions
|
||||
.get::<CodexFingerprintConvergenceContext>()
|
||||
.is_none()
|
||||
&& parts
|
||||
.extensions
|
||||
.get::<CodexFingerprintContextSlot>()
|
||||
.is_none()
|
||||
{
|
||||
parts
|
||||
.extensions
|
||||
.insert(CodexFingerprintContextSlot::default());
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn ensure_codex_fingerprint_context(
|
||||
parts: &mut Parts,
|
||||
body_json: &Value,
|
||||
) -> CodexFingerprintConvergenceContext {
|
||||
let context = resolve_codex_fingerprint_context(parts, body_json);
|
||||
if parts
|
||||
.extensions
|
||||
.get::<CodexFingerprintConvergenceContext>()
|
||||
.is_none()
|
||||
{
|
||||
parts.extensions.remove::<CodexFingerprintContextSlot>();
|
||||
parts.extensions.insert(context.clone());
|
||||
}
|
||||
context
|
||||
}
|
||||
|
||||
pub(crate) fn attach_codex_logical_turn_context(
|
||||
parts: &mut Parts,
|
||||
body_json: &Value,
|
||||
logical_turn_id: &str,
|
||||
) -> CodexFingerprintConvergenceContext {
|
||||
let context =
|
||||
build_codex_fingerprint_context(&parts.headers, body_json, logical_turn_id.to_string());
|
||||
parts.extensions.remove::<CodexFingerprintContextSlot>();
|
||||
parts.extensions.insert(context.clone());
|
||||
context
|
||||
}
|
||||
|
||||
pub(crate) fn restore_codex_logical_turn_context(
|
||||
parts: &mut Parts,
|
||||
context: &CodexFingerprintConvergenceContext,
|
||||
) {
|
||||
parts.extensions.remove::<CodexFingerprintContextSlot>();
|
||||
parts.extensions.insert(context.clone());
|
||||
}
|
||||
|
||||
fn build_codex_fingerprint_context(
|
||||
headers: &HeaderMap,
|
||||
body_json: &Value,
|
||||
logical_turn_id: String,
|
||||
) -> CodexFingerprintConvergenceContext {
|
||||
let signals = codex_request_signals_from_request(headers, Some(body_json));
|
||||
let mut context =
|
||||
CodexFingerprintConvergenceContext::new(logical_turn_id, current_unix_millis());
|
||||
|
||||
if let Some(turn_id) = signals.turn_id {
|
||||
context = context.with_original_turn_id(turn_id);
|
||||
}
|
||||
if let Some(session_id) = signals.thread_id.or(signals.session_id) {
|
||||
context = context.with_original_client_session_id(session_id);
|
||||
}
|
||||
if let Some(prompt_cache_key) = signals.prompt_cache_key {
|
||||
context = context.with_original_prompt_cache_key(prompt_cache_key);
|
||||
}
|
||||
|
||||
context
|
||||
}
|
||||
|
||||
fn current_unix_millis() -> u64 {
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_millis()
|
||||
.try_into()
|
||||
.unwrap_or(u64::MAX)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use http::HeaderValue;
|
||||
use serde_json::json;
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn request_signals_are_captured_once_for_the_logical_turn() {
|
||||
let request = http::Request::builder()
|
||||
.header("thread-id", "header-thread")
|
||||
.body(())
|
||||
.expect("request should build");
|
||||
let (mut parts, _) = request.into_parts();
|
||||
let body = json!({
|
||||
"prompt_cache_key": "client-cache",
|
||||
"client_metadata": {
|
||||
"turn_id": "client-turn",
|
||||
"thread_id": "body-thread"
|
||||
}
|
||||
});
|
||||
|
||||
let context = attach_codex_logical_turn_context(&mut parts, &body, "logical-turn");
|
||||
|
||||
assert_eq!(context.logical_turn_id(), "logical-turn");
|
||||
assert_eq!(context.original_turn_id(), Some("client-turn"));
|
||||
assert_eq!(context.original_client_session_id(), Some("header-thread"));
|
||||
assert_eq!(context.original_prompt_cache_key(), Some("client-cache"));
|
||||
assert_eq!(
|
||||
parts.extensions.get::<CodexFingerprintConvergenceContext>(),
|
||||
Some(&context)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn restored_context_wins_over_retry_request_signals() {
|
||||
let original = CodexFingerprintConvergenceContext::new("logical-turn", 1234)
|
||||
.with_original_turn_id("original-turn")
|
||||
.with_original_client_session_id("original-thread")
|
||||
.with_original_prompt_cache_key("original-cache");
|
||||
let request = http::Request::builder()
|
||||
.body(())
|
||||
.expect("request should build");
|
||||
let (mut parts, _) = request.into_parts();
|
||||
parts
|
||||
.headers
|
||||
.insert("thread-id", HeaderValue::from_static("retry-thread"));
|
||||
restore_codex_logical_turn_context(&mut parts, &original);
|
||||
|
||||
let resolved = resolve_codex_fingerprint_context(
|
||||
&parts,
|
||||
&json!({
|
||||
"prompt_cache_key": "retry-cache",
|
||||
"client_metadata": {"turn_id": "retry-turn"}
|
||||
}),
|
||||
);
|
||||
|
||||
assert_eq!(resolved, original);
|
||||
assert_eq!(resolved.turn_started_at_unix_ms(), 1234);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn generated_context_is_persisted_for_http_replanning() {
|
||||
let request = http::Request::builder()
|
||||
.header("session-id", "client-session")
|
||||
.body(())
|
||||
.expect("request should build");
|
||||
let (mut parts, _) = request.into_parts();
|
||||
let body = json!({
|
||||
"prompt_cache_key": "client-cache",
|
||||
"client_metadata": {"turn_id": "client-turn"}
|
||||
});
|
||||
|
||||
let first = ensure_codex_fingerprint_context(&mut parts, &body);
|
||||
let second = resolve_codex_fingerprint_context(
|
||||
&parts,
|
||||
&json!({
|
||||
"prompt_cache_key": "retry-cache",
|
||||
"client_metadata": {"turn_id": "retry-turn"}
|
||||
}),
|
||||
);
|
||||
|
||||
assert_eq!(second, first);
|
||||
assert_eq!(second.original_turn_id(), Some("client-turn"));
|
||||
assert_eq!(second.original_prompt_cache_key(), Some("client-cache"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn installed_slot_reuses_context_across_cloned_parts() {
|
||||
let request = http::Request::builder()
|
||||
.body(())
|
||||
.expect("request should build");
|
||||
let (mut parts, _) = request.into_parts();
|
||||
install_codex_fingerprint_context_slot(&mut parts);
|
||||
let cloned_parts = parts.clone();
|
||||
|
||||
let first = resolve_codex_fingerprint_context(
|
||||
&parts,
|
||||
&json!({
|
||||
"prompt_cache_key": "first-cache",
|
||||
"client_metadata": {"turn_id": "first-turn"}
|
||||
}),
|
||||
);
|
||||
let second = resolve_codex_fingerprint_context(
|
||||
&cloned_parts,
|
||||
&json!({
|
||||
"prompt_cache_key": "second-cache",
|
||||
"client_metadata": {"turn_id": "second-turn"}
|
||||
}),
|
||||
);
|
||||
|
||||
assert_eq!(second, first);
|
||||
assert_eq!(second.original_turn_id(), Some("first-turn"));
|
||||
assert_eq!(second.original_prompt_cache_key(), Some("first-cache"));
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
mod adaptation;
|
||||
pub(crate) mod api;
|
||||
pub(crate) mod codex_context;
|
||||
mod finalize;
|
||||
mod planner;
|
||||
mod pure;
|
||||
|
||||
@@ -2,6 +2,7 @@ use std::collections::BTreeMap;
|
||||
use std::time::Duration;
|
||||
|
||||
use aether_ai_serving::{run_ai_authenticated_decision_input, AiAuthenticatedDecisionInputPort};
|
||||
use aether_provider_transport::CodexFingerprintConvergenceContext;
|
||||
use aether_routing_core::{
|
||||
rank_vector_for_candidate, CandidateKind, ResolvedRoutingPolicy, RoutingCandidateFacts,
|
||||
RoutingCandidateTrace, RoutingDecisionTrace, RoutingPoolExpansionTrace, RoutingRulePhase,
|
||||
@@ -55,7 +56,7 @@ pub(crate) struct LocalRequestedModelDecisionInput {
|
||||
pub(crate) client_surface: Option<ClientSurface>,
|
||||
pub(crate) gateway_credential_carrier: Option<GatewayCredentialCarrier>,
|
||||
pub(crate) client_session_affinity: Option<ClientSessionAffinity>,
|
||||
pub(crate) original_client_session_id: Option<String>,
|
||||
pub(crate) codex_fingerprint_context: Option<CodexFingerprintConvergenceContext>,
|
||||
pub(crate) routing_policy: Option<ResolvedRoutingPolicy>,
|
||||
pub(crate) routing_trace_seed: Option<RoutingDecisionTrace>,
|
||||
pub(crate) routing_context: Option<LocalRoutingRequestContext>,
|
||||
@@ -376,13 +377,25 @@ fn apply_codex_oauth_fingerprint_convergence_to_decision(
|
||||
else {
|
||||
return;
|
||||
};
|
||||
crate::ai_serving::transport::apply_codex_oauth_fingerprint_convergence(
|
||||
transport,
|
||||
provider_api_format,
|
||||
input.original_client_session_id.as_deref(),
|
||||
&mut decision.provider_request_headers,
|
||||
provider_request_body,
|
||||
);
|
||||
let Some(context) = input.codex_fingerprint_context.as_ref() else {
|
||||
return;
|
||||
};
|
||||
let applied =
|
||||
crate::ai_serving::transport::apply_codex_oauth_fingerprint_convergence_with_context(
|
||||
transport,
|
||||
provider_api_format,
|
||||
context,
|
||||
&mut decision.provider_request_headers,
|
||||
provider_request_body,
|
||||
);
|
||||
if applied {
|
||||
decision.prompt_cache_key = provider_request_body
|
||||
.get("prompt_cache_key")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.map(ToOwned::to_owned);
|
||||
}
|
||||
}
|
||||
|
||||
struct GatewayAuthenticatedDecisionInputPort<'a> {
|
||||
@@ -471,7 +484,7 @@ pub(crate) fn build_local_requested_model_decision_input(
|
||||
client_surface: None,
|
||||
gateway_credential_carrier: None,
|
||||
client_session_affinity: None,
|
||||
original_client_session_id: None,
|
||||
codex_fingerprint_context: None,
|
||||
routing_policy: None,
|
||||
routing_trace_seed: None,
|
||||
routing_context: None,
|
||||
@@ -486,7 +499,8 @@ pub(crate) async fn attach_routing_policy_to_local_requested_model_input(
|
||||
body_json: &Value,
|
||||
client_api_format: &str,
|
||||
) -> Result<(), GatewayError> {
|
||||
input.original_client_session_id = original_client_session_id_from_headers(&parts.headers);
|
||||
input.codex_fingerprint_context =
|
||||
Some(crate::ai_serving::codex_context::resolve_codex_fingerprint_context(parts, body_json));
|
||||
let explicit_group = routing_header_value_str(&parts.headers, ROUTING_GROUP_HEADER);
|
||||
let selected_group = match state.routing_group_read_repository() {
|
||||
Some(repository) => {
|
||||
@@ -737,12 +751,6 @@ pub(crate) async fn attach_routing_policy_to_local_requested_model_input(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn original_client_session_id_from_headers(headers: &HeaderMap) -> Option<String> {
|
||||
routing_header_value_str(headers, "session-id")
|
||||
.or_else(|| routing_header_value_str(headers, "session_id"))
|
||||
.or_else(|| routing_header_value_str(headers, "x-session-id"))
|
||||
}
|
||||
|
||||
fn try_attach_static_default_routing_policy_to_input(
|
||||
input: &mut LocalRequestedModelDecisionInput,
|
||||
parts: &http::request::Parts,
|
||||
@@ -1106,38 +1114,6 @@ mod tests {
|
||||
GatewayProviderTransportProvider,
|
||||
};
|
||||
|
||||
#[test]
|
||||
fn original_client_session_id_accepts_live_header_as_fallback() {
|
||||
let headers = HeaderMap::from_iter([(
|
||||
HeaderName::from_static("x-session-id"),
|
||||
HeaderValue::from_static("live-thread-1"),
|
||||
)]);
|
||||
|
||||
assert_eq!(
|
||||
original_client_session_id_from_headers(&headers).as_deref(),
|
||||
Some("live-thread-1")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn original_client_session_id_prefers_responses_headers_over_live_fallback() {
|
||||
let headers = HeaderMap::from_iter([
|
||||
(
|
||||
HeaderName::from_static("session-id"),
|
||||
HeaderValue::from_static("responses-session"),
|
||||
),
|
||||
(
|
||||
HeaderName::from_static("x-session-id"),
|
||||
HeaderValue::from_static("live-thread"),
|
||||
),
|
||||
]);
|
||||
|
||||
assert_eq!(
|
||||
original_client_session_id_from_headers(&headers).as_deref(),
|
||||
Some("responses-session")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn explicit_routing_selection_cache_key_is_principal_specific() {
|
||||
let first = routing_group_selection_cache_key(
|
||||
@@ -1350,7 +1326,7 @@ mod tests {
|
||||
client_surface: None,
|
||||
gateway_credential_carrier: None,
|
||||
client_session_affinity: None,
|
||||
original_client_session_id: None,
|
||||
codex_fingerprint_context: None,
|
||||
routing_policy: None,
|
||||
routing_trace_seed: None,
|
||||
model_directive_policy: Default::default(),
|
||||
@@ -1599,7 +1575,7 @@ mod tests {
|
||||
client_surface: None,
|
||||
gateway_credential_carrier: None,
|
||||
client_session_affinity: None,
|
||||
original_client_session_id: None,
|
||||
codex_fingerprint_context: None,
|
||||
routing_policy: None,
|
||||
routing_trace_seed: None,
|
||||
model_directive_policy: Default::default(),
|
||||
@@ -1669,7 +1645,7 @@ mod tests {
|
||||
client_surface: None,
|
||||
gateway_credential_carrier: None,
|
||||
client_session_affinity: None,
|
||||
original_client_session_id: None,
|
||||
codex_fingerprint_context: None,
|
||||
routing_policy: None,
|
||||
routing_trace_seed: None,
|
||||
routing_context: None,
|
||||
@@ -1754,7 +1730,13 @@ mod tests {
|
||||
});
|
||||
let mut with_mutation = sample_decision_input();
|
||||
for input in [&mut no_context, &mut empty_mutation, &mut with_mutation] {
|
||||
input.original_client_session_id = Some("client-session-1".to_string());
|
||||
input.codex_fingerprint_context = Some(
|
||||
CodexFingerprintConvergenceContext::new(
|
||||
uuid::Uuid::new_v4().to_string(),
|
||||
1_756_668_000_000,
|
||||
)
|
||||
.with_original_client_session_id("client-session-1".to_string()),
|
||||
);
|
||||
}
|
||||
|
||||
let mut stable_identity = None;
|
||||
@@ -1801,6 +1783,10 @@ mod tests {
|
||||
.provider_request_body
|
||||
.as_ref()
|
||||
.expect("request body");
|
||||
assert_eq!(
|
||||
decision.prompt_cache_key.as_deref(),
|
||||
body.get("prompt_cache_key").and_then(Value::as_str)
|
||||
);
|
||||
assert_eq!(
|
||||
body["prompt_cache_key"],
|
||||
"172c39e6-c0a0-5a70-8b63-e0f8e0d185a3"
|
||||
@@ -2005,6 +1991,7 @@ mod tests {
|
||||
|
||||
let body = decision.provider_request_body.as_ref().expect("body");
|
||||
assert!(body.get("prompt_cache_key").is_none());
|
||||
assert!(decision.prompt_cache_key.is_none());
|
||||
assert!(body.get("client_metadata").is_none());
|
||||
assert!(!decision.provider_request_headers.contains_key("session-id"));
|
||||
assert!(!decision.provider_request_headers.contains_key("thread-id"));
|
||||
|
||||
@@ -377,7 +377,7 @@ mod tests {
|
||||
client_surface: None,
|
||||
gateway_credential_carrier: None,
|
||||
client_session_affinity: None,
|
||||
original_client_session_id: None,
|
||||
codex_fingerprint_context: None,
|
||||
routing_policy: None,
|
||||
routing_trace_seed: None,
|
||||
routing_context: None,
|
||||
|
||||
@@ -2182,7 +2182,7 @@ mod tests {
|
||||
client_surface: None,
|
||||
gateway_credential_carrier: None,
|
||||
client_session_affinity: None,
|
||||
original_client_session_id: None,
|
||||
codex_fingerprint_context: None,
|
||||
routing_policy: None,
|
||||
routing_trace_seed: None,
|
||||
routing_context: None,
|
||||
|
||||
@@ -60,6 +60,7 @@ pub(crate) mod windsurf {
|
||||
|
||||
pub(crate) use aether_provider_transport::{
|
||||
append_transport_diagnostics_to_value, apply_codex_oauth_fingerprint_convergence,
|
||||
apply_codex_oauth_fingerprint_convergence_with_context,
|
||||
apply_local_auth_config_header_overrides, apply_local_body_rules,
|
||||
apply_local_body_rules_with_request_headers, apply_local_header_rules,
|
||||
apply_local_header_rules_with_request_headers, apply_standard_provider_request_body_rules,
|
||||
|
||||
@@ -23,6 +23,21 @@ pub(crate) struct ClientSessionScope {
|
||||
pub(crate) source: ClientSessionSignalSource,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
pub(crate) struct CodexRequestSignals {
|
||||
pub(crate) session_id: Option<String>,
|
||||
pub(crate) thread_id: Option<String>,
|
||||
pub(crate) turn_id: Option<String>,
|
||||
pub(crate) prompt_cache_key: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
struct CodexTurnMetadataSignals {
|
||||
session_id: Option<String>,
|
||||
thread_id: Option<String>,
|
||||
turn_id: Option<String>,
|
||||
}
|
||||
|
||||
impl ClientSessionScope {
|
||||
fn new(
|
||||
client_family: impl Into<String>,
|
||||
@@ -136,6 +151,13 @@ pub(crate) fn client_session_scope_from_request(
|
||||
.or_else(|| extract_scope_from_other_specific_adapters(&request, client_family.as_str()))
|
||||
}
|
||||
|
||||
pub(crate) fn codex_request_signals_from_request(
|
||||
headers: &http::HeaderMap,
|
||||
body_json: Option<&Value>,
|
||||
) -> CodexRequestSignals {
|
||||
extract_codex_request_signals(&ClientSessionRequest { headers, body_json })
|
||||
}
|
||||
|
||||
fn codex_search_session_scope(request: &ClientSessionRequest<'_>) -> Option<ClientSessionScope> {
|
||||
let session_id = request
|
||||
.body_json?
|
||||
@@ -408,29 +430,7 @@ impl ClientSessionScopeAdapter for CodexSessionScopeAdapter {
|
||||
}
|
||||
|
||||
fn extract_scope(&self, request: &ClientSessionRequest<'_>) -> Option<ClientSessionScope> {
|
||||
header_value_str(request.headers, "session-id")
|
||||
.or_else(|| header_value_str(request.headers, "thread-id"))
|
||||
.or_else(|| header_value_str(request.headers, "session_id"))
|
||||
.or_else(|| header_value_str(request.headers, "conversation_id"))
|
||||
.map(|root_session| {
|
||||
ClientSessionScope::new(
|
||||
self.family(),
|
||||
root_session,
|
||||
None,
|
||||
header_value_str(request.headers, "chatgpt-account-id"),
|
||||
ClientSessionSignalSource::Header,
|
||||
)
|
||||
})
|
||||
.or_else(|| {
|
||||
let body_session = GenericSessionScopeAdapter.extract_scope(request)?;
|
||||
Some(ClientSessionScope::new(
|
||||
self.family(),
|
||||
body_session.session_id,
|
||||
body_session.agent_id,
|
||||
header_value_str(request.headers, "chatgpt-account-id"),
|
||||
body_session.source,
|
||||
))
|
||||
})
|
||||
codex_request_session_scope_from_request(request)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -778,6 +778,162 @@ fn explicit_aether_session_scope(
|
||||
))
|
||||
}
|
||||
|
||||
fn extract_codex_request_signals(request: &ClientSessionRequest<'_>) -> CodexRequestSignals {
|
||||
let body_client_metadata = request
|
||||
.body_json
|
||||
.and_then(|body| body.get("client_metadata"))
|
||||
.and_then(Value::as_object);
|
||||
let body_turn_metadata = codex_turn_metadata_signals(
|
||||
body_client_metadata.and_then(|metadata| metadata.get("x-codex-turn-metadata")),
|
||||
);
|
||||
let header_turn_metadata = header_value_str(request.headers, "x-codex-turn-metadata")
|
||||
.map(|raw| parse_codex_turn_metadata(&raw))
|
||||
.unwrap_or_default();
|
||||
|
||||
let native_thread_id = header_value_str(request.headers, "thread-id")
|
||||
.or_else(|| {
|
||||
body_client_metadata
|
||||
.and_then(|metadata| value_at_map_path(metadata, "thread_id"))
|
||||
.map(ToOwned::to_owned)
|
||||
})
|
||||
.or_else(|| body_turn_metadata.thread_id.clone());
|
||||
let turn_id = body_client_metadata
|
||||
.and_then(|metadata| value_at_map_path(metadata, "turn_id"))
|
||||
.map(ToOwned::to_owned)
|
||||
.or_else(|| body_turn_metadata.turn_id.clone())
|
||||
.or_else(|| {
|
||||
request
|
||||
.body_json
|
||||
.and_then(|body| value_at_path(body, &["turn_id"]))
|
||||
.map(ToOwned::to_owned)
|
||||
})
|
||||
.or(header_turn_metadata.turn_id);
|
||||
let prompt_cache_key = request
|
||||
.body_json
|
||||
.and_then(|body| value_at_path(body, &["prompt_cache_key"]))
|
||||
.map(ToOwned::to_owned);
|
||||
let session_id =
|
||||
codex_request_session_scope(request, &body_turn_metadata).map(|scope| scope.session_id);
|
||||
let thread_id = native_thread_id.or_else(|| session_id.clone());
|
||||
|
||||
CodexRequestSignals {
|
||||
session_id,
|
||||
thread_id,
|
||||
turn_id,
|
||||
prompt_cache_key,
|
||||
}
|
||||
}
|
||||
|
||||
fn codex_request_session_scope_from_request(
|
||||
request: &ClientSessionRequest<'_>,
|
||||
) -> Option<ClientSessionScope> {
|
||||
let body_turn_metadata = codex_turn_metadata_signals(
|
||||
request
|
||||
.body_json
|
||||
.and_then(|body| body.get("client_metadata"))
|
||||
.and_then(Value::as_object)
|
||||
.and_then(|metadata| metadata.get("x-codex-turn-metadata")),
|
||||
);
|
||||
codex_request_session_scope(request, &body_turn_metadata)
|
||||
}
|
||||
|
||||
fn codex_request_session_scope(
|
||||
request: &ClientSessionRequest<'_>,
|
||||
body_turn_metadata: &CodexTurnMetadataSignals,
|
||||
) -> Option<ClientSessionScope> {
|
||||
if let Some(scope) = explicit_aether_session_scope(request, CodexSessionScopeAdapter.family()) {
|
||||
return Some(scope);
|
||||
}
|
||||
|
||||
if let Some(root_session) = header_value_str(request.headers, "session-id")
|
||||
.or_else(|| header_value_str(request.headers, "thread-id"))
|
||||
.or_else(|| header_value_str(request.headers, "session_id"))
|
||||
.or_else(|| header_value_str(request.headers, "conversation_id"))
|
||||
.or_else(|| header_value_str(request.headers, "x-session-id"))
|
||||
{
|
||||
return Some(codex_session_scope(
|
||||
request,
|
||||
root_session,
|
||||
None,
|
||||
ClientSessionSignalSource::Header,
|
||||
));
|
||||
}
|
||||
|
||||
let body_client_metadata = request
|
||||
.body_json
|
||||
.and_then(|body| body.get("client_metadata"))
|
||||
.and_then(Value::as_object);
|
||||
if let Some(root_session) = body_client_metadata
|
||||
.and_then(|metadata| value_at_map_path(metadata, "session_id"))
|
||||
.or_else(|| {
|
||||
body_client_metadata.and_then(|metadata| value_at_map_path(metadata, "thread_id"))
|
||||
})
|
||||
.map(ToOwned::to_owned)
|
||||
.or_else(|| body_turn_metadata.session_id.clone())
|
||||
.or_else(|| body_turn_metadata.thread_id.clone())
|
||||
{
|
||||
return Some(codex_session_scope(
|
||||
request,
|
||||
root_session,
|
||||
None,
|
||||
ClientSessionSignalSource::Body,
|
||||
));
|
||||
}
|
||||
|
||||
let generic = GenericSessionScopeAdapter.extract_scope(request)?;
|
||||
Some(codex_session_scope(
|
||||
request,
|
||||
generic.session_id,
|
||||
generic.agent_id,
|
||||
generic.source,
|
||||
))
|
||||
}
|
||||
|
||||
fn codex_session_scope(
|
||||
request: &ClientSessionRequest<'_>,
|
||||
session_id: String,
|
||||
agent_id: Option<String>,
|
||||
source: ClientSessionSignalSource,
|
||||
) -> ClientSessionScope {
|
||||
ClientSessionScope::new(
|
||||
CodexSessionScopeAdapter.family(),
|
||||
session_id,
|
||||
agent_id,
|
||||
header_value_str(request.headers, "chatgpt-account-id"),
|
||||
source,
|
||||
)
|
||||
}
|
||||
|
||||
fn codex_turn_metadata_signals(value: Option<&Value>) -> CodexTurnMetadataSignals {
|
||||
match value {
|
||||
Some(Value::Object(metadata)) => codex_turn_metadata_signals_from_map(metadata),
|
||||
Some(Value::String(raw)) => parse_codex_turn_metadata(raw),
|
||||
_ => CodexTurnMetadataSignals::default(),
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_codex_turn_metadata(raw: &str) -> CodexTurnMetadataSignals {
|
||||
serde_json::from_str::<Map<String, Value>>(raw)
|
||||
.map(|metadata| codex_turn_metadata_signals_from_map(&metadata))
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
fn codex_turn_metadata_signals_from_map(metadata: &Map<String, Value>) -> CodexTurnMetadataSignals {
|
||||
CodexTurnMetadataSignals {
|
||||
session_id: value_at_map_path(metadata, "session_id").map(ToOwned::to_owned),
|
||||
thread_id: value_at_map_path(metadata, "thread_id").map(ToOwned::to_owned),
|
||||
turn_id: value_at_map_path(metadata, "turn_id").map(ToOwned::to_owned),
|
||||
}
|
||||
}
|
||||
|
||||
fn value_at_map_path<'a>(object: &'a Map<String, Value>, key: &str) -> Option<&'a str> {
|
||||
object
|
||||
.get(key)
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
}
|
||||
|
||||
fn normalize_session_key(
|
||||
account_hint: Option<&str>,
|
||||
root_session: &str,
|
||||
@@ -831,12 +987,25 @@ mod tests {
|
||||
client_session_affinity_from_api_request,
|
||||
client_session_affinity_from_report_context_value, client_session_affinity_from_request,
|
||||
client_session_affinity_report_context_value, client_session_scope_from_request,
|
||||
ClientSessionSignalSource, AETHER_AGENT_ID_HEADER, AETHER_SESSION_ID_HEADER,
|
||||
codex_request_signals_from_request, ClientSessionSignalSource, AETHER_AGENT_ID_HEADER,
|
||||
AETHER_SESSION_ID_HEADER,
|
||||
};
|
||||
use aether_scheduler_core::ClientSessionAffinity;
|
||||
use http::{HeaderMap, HeaderValue};
|
||||
use http::{HeaderMap, HeaderName, HeaderValue};
|
||||
use serde_json::json;
|
||||
|
||||
fn request_headers(values: &[(&str, &str)]) -> HeaderMap {
|
||||
values
|
||||
.iter()
|
||||
.map(|(name, value)| {
|
||||
(
|
||||
HeaderName::from_bytes(name.as_bytes()).expect("valid test header name"),
|
||||
HeaderValue::from_bytes(value.as_bytes()).expect("valid test header value"),
|
||||
)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_adapter_extracts_body_session_and_agent() {
|
||||
let body = json!({
|
||||
@@ -933,6 +1102,276 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_request_signals_apply_session_precedence() {
|
||||
let cases = vec![
|
||||
(
|
||||
request_headers(&[
|
||||
(AETHER_SESSION_ID_HEADER, "aether-session"),
|
||||
("session-id", "header-session"),
|
||||
]),
|
||||
json!({"client_metadata": {"session_id": "body-session"}}),
|
||||
"aether-session",
|
||||
ClientSessionSignalSource::ExplicitAetherHeader,
|
||||
),
|
||||
(
|
||||
request_headers(&[
|
||||
("session-id", "header-session"),
|
||||
("thread-id", "header-thread"),
|
||||
("session_id", "header-session-underscore"),
|
||||
("conversation_id", "header-conversation"),
|
||||
]),
|
||||
json!({"client_metadata": {"session_id": "body-session"}}),
|
||||
"header-session",
|
||||
ClientSessionSignalSource::Header,
|
||||
),
|
||||
(
|
||||
request_headers(&[
|
||||
("thread-id", "header-thread"),
|
||||
("session_id", "header-session-underscore"),
|
||||
("conversation_id", "header-conversation"),
|
||||
]),
|
||||
json!({"client_metadata": {"session_id": "body-session"}}),
|
||||
"header-thread",
|
||||
ClientSessionSignalSource::Header,
|
||||
),
|
||||
(
|
||||
request_headers(&[
|
||||
("session_id", "header-session-underscore"),
|
||||
("conversation_id", "header-conversation"),
|
||||
]),
|
||||
json!({"client_metadata": {"session_id": "body-session"}}),
|
||||
"header-session-underscore",
|
||||
ClientSessionSignalSource::Header,
|
||||
),
|
||||
(
|
||||
request_headers(&[("conversation_id", "header-conversation")]),
|
||||
json!({"client_metadata": {"session_id": "body-session"}}),
|
||||
"header-conversation",
|
||||
ClientSessionSignalSource::Header,
|
||||
),
|
||||
(
|
||||
HeaderMap::new(),
|
||||
json!({
|
||||
"prompt_cache_key": "prompt-cache",
|
||||
"client_metadata": {
|
||||
"session_id": "body-session",
|
||||
"thread_id": "body-thread",
|
||||
"x-codex-turn-metadata": {
|
||||
"session_id": "nested-session",
|
||||
"thread_id": "nested-thread"
|
||||
}
|
||||
}
|
||||
}),
|
||||
"body-session",
|
||||
ClientSessionSignalSource::Body,
|
||||
),
|
||||
(
|
||||
HeaderMap::new(),
|
||||
json!({
|
||||
"prompt_cache_key": "prompt-cache",
|
||||
"client_metadata": {
|
||||
"thread_id": "body-thread",
|
||||
"x-codex-turn-metadata": {"session_id": "nested-session"}
|
||||
}
|
||||
}),
|
||||
"body-thread",
|
||||
ClientSessionSignalSource::Body,
|
||||
),
|
||||
(
|
||||
HeaderMap::new(),
|
||||
json!({
|
||||
"prompt_cache_key": "prompt-cache",
|
||||
"client_metadata": {
|
||||
"x-codex-turn-metadata": {
|
||||
"session_id": "nested-session",
|
||||
"thread_id": "nested-thread"
|
||||
}
|
||||
}
|
||||
}),
|
||||
"nested-session",
|
||||
ClientSessionSignalSource::Body,
|
||||
),
|
||||
(
|
||||
HeaderMap::new(),
|
||||
json!({
|
||||
"prompt_cache_key": "prompt-cache",
|
||||
"client_metadata": {
|
||||
"x-codex-turn-metadata": json!({
|
||||
"thread_id": "nested-thread"
|
||||
}).to_string()
|
||||
}
|
||||
}),
|
||||
"nested-thread",
|
||||
ClientSessionSignalSource::Body,
|
||||
),
|
||||
(
|
||||
HeaderMap::new(),
|
||||
json!({
|
||||
"prompt_cache_key": "prompt-cache",
|
||||
"conversation_id": "generic-conversation"
|
||||
}),
|
||||
"prompt-cache",
|
||||
ClientSessionSignalSource::Body,
|
||||
),
|
||||
(
|
||||
HeaderMap::new(),
|
||||
json!({"metadata": {"session_id": "generic-session"}}),
|
||||
"generic-session",
|
||||
ClientSessionSignalSource::Body,
|
||||
),
|
||||
];
|
||||
|
||||
for (headers, body, expected_session_id, expected_source) in cases {
|
||||
let signals = codex_request_signals_from_request(&headers, Some(&body));
|
||||
assert_eq!(signals.session_id.as_deref(), Some(expected_session_id));
|
||||
|
||||
let mut codex_headers = headers;
|
||||
codex_headers.insert(
|
||||
http::header::USER_AGENT,
|
||||
HeaderValue::from_static("codex_cli_rs/0.144.1"),
|
||||
);
|
||||
let scope = client_session_scope_from_request(&codex_headers, Some(&body))
|
||||
.expect("Codex scope should reuse the native signal precedence");
|
||||
assert_eq!(scope.client_family, "codex");
|
||||
assert_eq!(scope.session_id, expected_session_id);
|
||||
assert_eq!(scope.source, expected_source);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_request_signals_extract_thread_and_prompt_cache_independently() {
|
||||
let body = json!({
|
||||
"prompt_cache_key": "prompt-cache",
|
||||
"client_metadata": {
|
||||
"thread_id": "body-thread",
|
||||
"x-codex-turn-metadata": {"thread_id": "nested-thread"}
|
||||
}
|
||||
});
|
||||
let headers = request_headers(&[("thread-id", "header-thread")]);
|
||||
let header_signals = codex_request_signals_from_request(&headers, Some(&body));
|
||||
assert_eq!(header_signals.thread_id.as_deref(), Some("header-thread"));
|
||||
assert_eq!(
|
||||
header_signals.prompt_cache_key.as_deref(),
|
||||
Some("prompt-cache")
|
||||
);
|
||||
|
||||
let body_signals = codex_request_signals_from_request(&HeaderMap::new(), Some(&body));
|
||||
assert_eq!(body_signals.thread_id.as_deref(), Some("body-thread"));
|
||||
|
||||
let nested_body = json!({
|
||||
"client_metadata": {
|
||||
"x-codex-turn-metadata": json!({
|
||||
"thread_id": "nested-thread"
|
||||
}).to_string()
|
||||
}
|
||||
});
|
||||
let nested_signals =
|
||||
codex_request_signals_from_request(&HeaderMap::new(), Some(&nested_body));
|
||||
assert_eq!(nested_signals.thread_id.as_deref(), Some("nested-thread"));
|
||||
|
||||
let session_only_body = json!({"client_metadata": {"session_id": "body-session"}});
|
||||
let session_only_signals =
|
||||
codex_request_signals_from_request(&HeaderMap::new(), Some(&session_only_body));
|
||||
assert_eq!(
|
||||
session_only_signals.thread_id.as_deref(),
|
||||
Some("body-session")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_request_signals_use_live_session_header() {
|
||||
let headers = request_headers(&[("x-session-id", "live-session")]);
|
||||
let signals = codex_request_signals_from_request(&headers, None);
|
||||
|
||||
assert_eq!(signals.session_id.as_deref(), Some("live-session"));
|
||||
assert_eq!(signals.thread_id.as_deref(), Some("live-session"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_request_signals_prefer_responses_session_header_over_live_session_header() {
|
||||
let headers = request_headers(&[
|
||||
("session-id", "responses-session"),
|
||||
("x-session-id", "live-session"),
|
||||
]);
|
||||
let signals = codex_request_signals_from_request(&headers, None);
|
||||
|
||||
assert_eq!(signals.session_id.as_deref(), Some("responses-session"));
|
||||
assert_eq!(signals.thread_id.as_deref(), Some("responses-session"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_request_signals_apply_turn_precedence() {
|
||||
let headers = request_headers(&[("x-codex-turn-metadata", r#"{"turn_id":"header-turn"}"#)]);
|
||||
let direct_body = json!({
|
||||
"turn_id": "top-level-turn",
|
||||
"client_metadata": {
|
||||
"turn_id": "body-turn",
|
||||
"x-codex-turn-metadata": {"turn_id": "nested-turn"}
|
||||
}
|
||||
});
|
||||
assert_eq!(
|
||||
codex_request_signals_from_request(&headers, Some(&direct_body))
|
||||
.turn_id
|
||||
.as_deref(),
|
||||
Some("body-turn")
|
||||
);
|
||||
|
||||
let nested_object_body = json!({
|
||||
"turn_id": "top-level-turn",
|
||||
"client_metadata": {
|
||||
"x-codex-turn-metadata": {"turn_id": "nested-object-turn"}
|
||||
}
|
||||
});
|
||||
assert_eq!(
|
||||
codex_request_signals_from_request(&headers, Some(&nested_object_body))
|
||||
.turn_id
|
||||
.as_deref(),
|
||||
Some("nested-object-turn")
|
||||
);
|
||||
|
||||
let nested_string_body = json!({
|
||||
"turn_id": "top-level-turn",
|
||||
"client_metadata": {
|
||||
"x-codex-turn-metadata": json!({
|
||||
"turn_id": "nested-string-turn"
|
||||
}).to_string()
|
||||
}
|
||||
});
|
||||
assert_eq!(
|
||||
codex_request_signals_from_request(&headers, Some(&nested_string_body))
|
||||
.turn_id
|
||||
.as_deref(),
|
||||
Some("nested-string-turn")
|
||||
);
|
||||
|
||||
let top_level_body = json!({
|
||||
"turn_id": "top-level-turn",
|
||||
"client_metadata": {"x-codex-turn-metadata": "not-json"}
|
||||
});
|
||||
assert_eq!(
|
||||
codex_request_signals_from_request(&headers, Some(&top_level_body))
|
||||
.turn_id
|
||||
.as_deref(),
|
||||
Some("top-level-turn")
|
||||
);
|
||||
assert_eq!(
|
||||
codex_request_signals_from_request(&headers, None)
|
||||
.turn_id
|
||||
.as_deref(),
|
||||
Some("header-turn")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn codex_request_signals_ignore_client_request_id() {
|
||||
let headers = request_headers(&[("x-client-request-id", "request-only-id")]);
|
||||
let signals =
|
||||
codex_request_signals_from_request(&headers, Some(&json!({"model": "gpt-5"})));
|
||||
|
||||
assert_eq!(signals, super::CodexRequestSignals::default());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn report_context_round_trips_normalized_session_affinity() {
|
||||
let affinity = ClientSessionAffinity::new(
|
||||
|
||||
@@ -1092,6 +1092,7 @@ async fn proxy_request_inner(
|
||||
),
|
||||
}
|
||||
let (mut parts, body) = request.into_parts();
|
||||
crate::ai_serving::codex_context::install_codex_fingerprint_context_slot(&mut parts);
|
||||
let redaction_slot = crate::privacy::RedactionSessionSlot::default();
|
||||
parts.extensions.insert(redaction_slot.clone());
|
||||
parts
|
||||
|
||||
@@ -49,6 +49,8 @@ pub(super) enum LiveAuthMode {
|
||||
pub(super) struct PlannedLiveCandidate {
|
||||
pub(super) execution: AiExecutionDecision,
|
||||
pub(super) pinned_candidate: ResponsesWebSocketPinnedCandidate,
|
||||
pub(super) codex_fingerprint_context:
|
||||
aether_provider_transport::CodexFingerprintConvergenceContext,
|
||||
pub(super) client_model: String,
|
||||
pub(super) provider_model: String,
|
||||
pub(super) auth_mode: LiveAuthMode,
|
||||
@@ -228,8 +230,9 @@ async fn plan_live_candidate_inner(
|
||||
if validate_model(client_model).is_err() || client_model.len() > MAX_LIVE_MODEL_BYTES {
|
||||
return Ok(None);
|
||||
}
|
||||
let parts = build_live_planning_parts(headers, remote_addr);
|
||||
let mut parts = build_live_planning_parts(headers, remote_addr);
|
||||
let body = json!({"model": client_model, "input": []});
|
||||
crate::ai_serving::codex_context::install_codex_fingerprint_context_slot(&mut parts);
|
||||
let execution = maybe_build_pinned_stream_local_same_format_provider_decision_payload(
|
||||
state,
|
||||
&parts,
|
||||
@@ -338,6 +341,8 @@ async fn plan_live_candidate_inner(
|
||||
Ok(Some(PlannedLiveCandidate {
|
||||
execution,
|
||||
pinned_candidate,
|
||||
codex_fingerprint_context:
|
||||
crate::ai_serving::codex_context::resolve_codex_fingerprint_context(&parts, &body),
|
||||
client_model: client_model.to_string(),
|
||||
provider_model,
|
||||
auth_mode,
|
||||
@@ -560,7 +565,11 @@ pub(super) fn build_live_stream_admission_attempt(
|
||||
remote_addr: &SocketAddr,
|
||||
upstream_url: String,
|
||||
) -> Result<Option<AiStreamAttempt>, GatewayError> {
|
||||
let parts = build_live_planning_parts(headers, remote_addr);
|
||||
let mut parts = build_live_planning_parts(headers, remote_addr);
|
||||
crate::ai_serving::codex_context::restore_codex_logical_turn_context(
|
||||
&mut parts,
|
||||
&candidate.codex_fingerprint_context,
|
||||
);
|
||||
let body = json!({"model": candidate.client_model.as_str(), "input": []});
|
||||
let mut execution = candidate.execution.clone();
|
||||
execution.upstream_url = Some(upstream_url);
|
||||
@@ -922,6 +931,11 @@ mod tests {
|
||||
"key-1",
|
||||
)
|
||||
.unwrap(),
|
||||
codex_fingerprint_context:
|
||||
aether_provider_transport::CodexFingerprintConvergenceContext::new(
|
||||
"test-live-turn",
|
||||
1,
|
||||
),
|
||||
client_model: "global-model".to_string(),
|
||||
provider_model: "provider-model".to_string(),
|
||||
auth_mode,
|
||||
|
||||
@@ -718,6 +718,11 @@ mod tests {
|
||||
PlannedLiveCandidate {
|
||||
execution,
|
||||
pinned_candidate: binding.pinned_candidate.clone(),
|
||||
codex_fingerprint_context:
|
||||
aether_provider_transport::CodexFingerprintConvergenceContext::new(
|
||||
"test-live-turn",
|
||||
1,
|
||||
),
|
||||
client_model: binding.client_model.clone(),
|
||||
provider_model: binding.provider_model.clone(),
|
||||
auth_mode: binding.auth_mode,
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
//! Client-side Responses WebSocket event forwarding and follow-up planning.
|
||||
|
||||
use aether_provider_transport::CodexFingerprintConvergenceContext;
|
||||
use axum::extract::ws::{Message as AxumWsMessage, WebSocket};
|
||||
use futures_util::SinkExt;
|
||||
use serde_json::Value;
|
||||
@@ -248,7 +249,14 @@ pub(super) async fn forward_client_message(
|
||||
// derive one strong live control snapshot that every stage below
|
||||
// shares. The connection's Upgrade-time decision is only the
|
||||
// immutable identity seed.
|
||||
let planning_parts = build_planning_parts(context);
|
||||
let logical_turn_id = Uuid::now_v7().to_string();
|
||||
let mut planning_parts = build_planning_parts(context);
|
||||
let codex_fingerprint_context =
|
||||
crate::ai_serving::codex_context::attach_codex_logical_turn_context(
|
||||
&mut planning_parts,
|
||||
&client_event,
|
||||
&logical_turn_id,
|
||||
);
|
||||
let turn_control = match resolve_responses_websocket_turn_control(
|
||||
state,
|
||||
context,
|
||||
@@ -453,6 +461,8 @@ pub(super) async fn forward_client_message(
|
||||
context,
|
||||
planning_parts,
|
||||
client_event,
|
||||
logical_turn_id,
|
||||
codex_fingerprint_context,
|
||||
turn_control,
|
||||
turn_redaction_session,
|
||||
)
|
||||
@@ -466,6 +476,8 @@ pub(super) async fn forward_client_message(
|
||||
planning_parts,
|
||||
client_event,
|
||||
requested_model,
|
||||
logical_turn_id,
|
||||
codex_fingerprint_context,
|
||||
turn_control,
|
||||
raw_responses_lite_static_config
|
||||
.expect("independent turns always retain their raw static config"),
|
||||
@@ -513,6 +525,8 @@ async fn forward_pinned_continuation(
|
||||
context: &WebSocketRequestContext,
|
||||
planning_parts: http::request::Parts,
|
||||
client_event: Value,
|
||||
logical_turn_id: String,
|
||||
codex_fingerprint_context: CodexFingerprintConvergenceContext,
|
||||
turn_control: ResponsesWebSocketTurnControl,
|
||||
turn_redaction_session: Option<RedactionSession>,
|
||||
) -> RelayDisposition {
|
||||
@@ -556,7 +570,6 @@ async fn forward_pinned_continuation(
|
||||
};
|
||||
|
||||
let turn_request_id = Uuid::new_v4().to_string();
|
||||
let logical_turn_id = Uuid::new_v4().to_string();
|
||||
let planned = match await_owned_responses_websocket_plan(spawn_owned_responses_websocket_plan(
|
||||
state.clone(),
|
||||
planning_parts,
|
||||
@@ -766,6 +779,7 @@ async fn forward_pinned_continuation(
|
||||
bound.body_normalization = normalization;
|
||||
bound.turn_state.begin(
|
||||
LogicalTurn::new(client_event, turn_index, logical_turn_id)
|
||||
.with_codex_fingerprint_context(codex_fingerprint_context)
|
||||
.with_provider_store(provider_event.get("store") == Some(&Value::Bool(true)))
|
||||
.with_turn_control(turn_control),
|
||||
turn,
|
||||
@@ -800,12 +814,13 @@ async fn forward_replanned_response_create(
|
||||
planning_parts: http::request::Parts,
|
||||
client_event: Value,
|
||||
requested_model: String,
|
||||
logical_turn_id: String,
|
||||
codex_fingerprint_context: CodexFingerprintConvergenceContext,
|
||||
turn_control: ResponsesWebSocketTurnControl,
|
||||
raw_responses_lite_static_config: ResponsesLiteStaticConfig,
|
||||
turn_redaction_session: Option<RedactionSession>,
|
||||
) -> RelayDisposition {
|
||||
let turn_request_id = Uuid::new_v4().to_string();
|
||||
let logical_turn_id = Uuid::new_v4().to_string();
|
||||
let now_unix_secs = current_unix_secs();
|
||||
let excluded_key_ids = bound.exhausted_exclusions.key_ids(now_unix_secs);
|
||||
let excluded_codex_account_ids = bound.exhausted_exclusions.codex_account_ids(now_unix_secs);
|
||||
@@ -992,6 +1007,7 @@ async fn forward_replanned_response_create(
|
||||
bound.body_normalization = normalization;
|
||||
bound.turn_state.begin(
|
||||
LogicalTurn::new(client_event.clone(), turn_index, logical_turn_id.clone())
|
||||
.with_codex_fingerprint_context(codex_fingerprint_context.clone())
|
||||
.with_provider_store(provider_event.get("store") == Some(&Value::Bool(true)))
|
||||
.with_turn_control(turn_control),
|
||||
turn,
|
||||
@@ -1076,6 +1092,7 @@ async fn forward_replanned_response_create(
|
||||
bound.binding_identity = replacement.binding_identity;
|
||||
bound.turn_state.begin(
|
||||
LogicalTurn::new(client_event, turn_index, logical_turn_id)
|
||||
.with_codex_fingerprint_context(codex_fingerprint_context)
|
||||
.with_provider_store(provider_event.get("store") == Some(&Value::Bool(true)))
|
||||
.with_turn_control(turn_control),
|
||||
turn,
|
||||
|
||||
@@ -146,6 +146,7 @@ pub(super) async fn retry_active_turn_after_quota_exhaustion(
|
||||
};
|
||||
let turn_index = active.turn_index;
|
||||
let logical_turn_id = active.logical_turn_id.clone();
|
||||
let codex_fingerprint_context = active.codex_fingerprint_context.clone();
|
||||
let turn_attempt = active.turn_attempt;
|
||||
|
||||
let retry_exclusion_until_unix_secs = bound
|
||||
@@ -154,7 +155,13 @@ pub(super) async fn retry_active_turn_after_quota_exhaustion(
|
||||
let exhausted_key = record_exhausted_bound_key(bound, retry_exclusion_until_unix_secs);
|
||||
let exhausted_key_id = exhausted_key.as_ref().map(|(key_id, _)| key_id.clone());
|
||||
|
||||
let planning_parts = build_planning_parts(context);
|
||||
let mut planning_parts = build_planning_parts(context);
|
||||
if let Some(codex_fingerprint_context) = codex_fingerprint_context.as_ref() {
|
||||
crate::ai_serving::codex_context::restore_codex_logical_turn_context(
|
||||
&mut planning_parts,
|
||||
codex_fingerprint_context,
|
||||
);
|
||||
}
|
||||
let turn_request_id = Uuid::new_v4().to_string();
|
||||
let now_unix_secs = current_unix_secs();
|
||||
let excluded_key_ids = bound.exhausted_exclusions.key_ids(now_unix_secs);
|
||||
|
||||
@@ -444,7 +444,14 @@ async fn bootstrap_responses_websocket(
|
||||
let raw_responses_lite_static_config =
|
||||
ResponsesLiteStaticConfig::from_response_create(&first_event);
|
||||
|
||||
let planning_parts = build_planning_parts(context);
|
||||
let first_logical_turn_id = Uuid::now_v7().to_string();
|
||||
let mut planning_parts = build_planning_parts(context);
|
||||
let first_codex_fingerprint_context =
|
||||
crate::ai_serving::codex_context::attach_codex_logical_turn_context(
|
||||
&mut planning_parts,
|
||||
&first_event,
|
||||
&first_logical_turn_id,
|
||||
);
|
||||
let turn_control = match resolve_responses_websocket_turn_control(
|
||||
&state,
|
||||
context,
|
||||
@@ -860,7 +867,6 @@ async fn bootstrap_responses_websocket(
|
||||
return None;
|
||||
}
|
||||
};
|
||||
let first_logical_turn_id = Uuid::new_v4().to_string();
|
||||
let first_turn_decision = prepare_responses_websocket_turn_decision(
|
||||
&decision,
|
||||
context.trace_id.clone(),
|
||||
@@ -954,6 +960,7 @@ async fn bootstrap_responses_websocket(
|
||||
}
|
||||
bound.turn_state.begin(
|
||||
LogicalTurn::new(first_event, 1, first_logical_turn_id)
|
||||
.with_codex_fingerprint_context(first_codex_fingerprint_context)
|
||||
.with_provider_store(first_provider_event.get("store") == Some(&Value::Bool(true)))
|
||||
.with_turn_control(turn_control),
|
||||
first_turn,
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
//! 非法组合只能靠调用点的 if 和「记得同时改另外两个字段」来避免。这里把它收敛成
|
||||
//! 一个枚举:合法组合由类型保证,转换只能走受控 API。
|
||||
|
||||
use aether_provider_transport::CodexFingerprintConvergenceContext;
|
||||
use serde_json::Value;
|
||||
|
||||
use super::control::ResponsesWebSocketTurnControl;
|
||||
@@ -26,6 +27,9 @@ pub(super) struct LogicalTurn {
|
||||
pub(super) provider_store: bool,
|
||||
pub(super) turn_index: u64,
|
||||
pub(super) logical_turn_id: String,
|
||||
/// Immutable Codex client identity for every provider attempt belonging to
|
||||
/// this logical turn. A transparent re-plan must never mint a new turn.
|
||||
pub(super) codex_fingerprint_context: Option<CodexFingerprintConvergenceContext>,
|
||||
pub(super) turn_attempt: u32,
|
||||
pub(super) retry_attempted: bool,
|
||||
pub(super) retry_unsafe_reason: Option<&'static str>,
|
||||
@@ -42,6 +46,7 @@ impl LogicalTurn {
|
||||
provider_store: false,
|
||||
turn_index,
|
||||
logical_turn_id,
|
||||
codex_fingerprint_context: None,
|
||||
turn_attempt: 1,
|
||||
retry_attempted: false,
|
||||
retry_unsafe_reason: None,
|
||||
@@ -54,6 +59,14 @@ impl LogicalTurn {
|
||||
self
|
||||
}
|
||||
|
||||
pub(super) fn with_codex_fingerprint_context(
|
||||
mut self,
|
||||
context: CodexFingerprintConvergenceContext,
|
||||
) -> Self {
|
||||
self.codex_fingerprint_context = Some(context);
|
||||
self
|
||||
}
|
||||
|
||||
pub(super) fn with_provider_store(mut self, provider_store: bool) -> Self {
|
||||
self.provider_store = provider_store;
|
||||
self
|
||||
|
||||
Reference in New Issue
Block a user