Skip to main content

kcode_intelligence_router/
lib.rs

1//! Typed, in-process intelligence operations for Kennedy.
2//!
3//! This crate has no HTTP router, health endpoint, provider catalog, or JSON
4//! request dispatch. Callers choose exact models and bind every provider call
5//! to a user before invoking it.
6
7#![forbid(unsafe_code)]
8
9mod defaults;
10mod error;
11mod pricing;
12mod receipts;
13
14use std::{
15    collections::{HashMap, VecDeque},
16    path::PathBuf,
17    sync::{Arc, Mutex},
18    time::{Duration, Instant},
19};
20
21use anyhow::Context;
22use chrono::{DateTime, NaiveDate, Utc};
23pub use kcode_codex_runtime::ReasoningEffort;
24use kcode_codex_runtime::{
25    CatalogCache, Codex, CodexConfig, ErrorKind as CodexErrorKind,
26    GenerationRequest as CodexGenerationRequest, TokenUsage as CodexTokenUsage,
27    WebSearchRequest as CodexSearchRequest,
28};
29use kcode_doc_extraction::{DocumentExtractor, DocumentInput, ErrorKind as DocumentErrorKind};
30use kcode_gemini_api::{
31    AgentTool as GeminiAgentTool, AgentTurnRequest as GeminiAgentRequest,
32    CompletionStatus as GeminiCompletionStatus, Error as GeminiError, Gemini, GenerationOptions,
33    GroundedSearchRequest, MediaInput as GeminiMediaInput, MultimodalRequest, NanoBananaProRequest,
34    ServiceTier, StructuredOutput, TextModel, ThinkingLevel, TokenUsage as GeminiTokenUsage,
35};
36use kcode_intelligence_cancellation::{ActiveOperation, OperationRegistry, RegistryError};
37use kcode_openai_api::{
38    AgentTool as OpenAiAgentTool, AgentTurnRequest as OpenAiAgentRequest, AudioInput,
39    Error as OpenAiError, ImageAnalysisRequest, ImageAnalysisStatus as OpenAiImageStatus,
40    ImageAnalysisUsage, ImageEditRequest as OpenAiImageEditRequest,
41    ImageGenerationRequest as OpenAiImageRequest, ImageInput as OpenAiImageInput,
42    ImageMediaType as OpenAiImageMediaType, ImageUsage as OpenAiGenerationUsage, OpenAi,
43    TranscriptionRequest as OpenAiTranscriptionRequest, TranscriptionUsage,
44};
45use kcode_web_fetch::{ErrorKind as WebFetchErrorKind, WebFetcher};
46use serde::{Deserialize, Serialize};
47use serde_json::{Value, json};
48use uuid::Uuid;
49
50use defaults::*;
51pub use error::{Error, ErrorKind, Result};
52pub use pricing::{
53    CostAccuracy, CostEstimate, PRICING_VERSION, estimate_cost, estimate_token_cost,
54};
55use receipts::ReceiptStore;
56pub use receipts::{DailyUsage, DailyUsageKey, Metering, TokenUsage, UsageReceipt};
57
58/// A successful provider result together with its canonical durable receipt.
59#[derive(Clone, Debug, PartialEq)]
60pub struct Accounted<T> {
61    /// Provider-normalized operation result.
62    pub value: T,
63    /// Canonical receipt written by the router for this call.
64    pub receipt: UsageReceipt,
65}
66
67impl<T> Accounted<T> {
68    /// Maps the operation value without separating it from its receipt.
69    pub fn map<U>(self, map: impl FnOnce(T) -> U) -> Accounted<U> {
70        Accounted {
71            value: map(self.value),
72            receipt: self.receipt,
73        }
74    }
75}
76
77/// Construction inputs for the intelligence library.
78pub struct Config {
79    pub openai_api_key: Option<String>,
80    pub gemini_api_key: Option<String>,
81    pub codex_catalog_cache: CatalogCache,
82    pub receipt_directory: PathBuf,
83}
84
85/// The configured Kennedy model, returned directly instead of through a catalog endpoint.
86#[derive(Clone, Debug, Eq, PartialEq)]
87pub struct RuntimeModel {
88    pub model: String,
89    pub reasoning_effort: String,
90    pub context_window_tokens: u64,
91    pub max_input_tokens: u64,
92}
93
94/// Provider selected for one resolved agent model.
95#[derive(Clone, Copy, Debug, Eq, PartialEq)]
96pub enum AgentProvider {
97    /// The locally authenticated Codex runtime.
98    Codex,
99    /// The direct OpenAI API.
100    OpenAi,
101    /// The direct Gemini API.
102    Gemini,
103}
104
105/// Exact provider model and capacity selected for an agent request.
106#[derive(Clone, Debug, Eq, PartialEq)]
107pub struct ResolvedAgentModel {
108    /// Caller-supplied model selector.
109    pub requested_model: String,
110    /// Exact provider model identifier.
111    pub provider_model: String,
112    /// Selected provider.
113    pub provider: AgentProvider,
114    /// Total context window.
115    pub context_window_tokens: u64,
116    /// Maximum input allowed by Kennedy.
117    pub max_input_tokens: u64,
118}
119
120#[derive(Clone)]
121pub struct Intelligence {
122    codex: Codex,
123    agent: kcode_codex_runtime_v2::Codex,
124    openai: Option<OpenAi>,
125    gemini: Option<Gemini>,
126    web_fetcher: WebFetcher,
127    document_extractor: DocumentExtractor,
128    active_operations: OperationRegistry,
129    receipts: ReceiptStore,
130    model_cache: Arc<Mutex<HashMap<String, ResolvedAgentModel>>>,
131}
132
133/// An intelligence handle that cannot make a provider call without user attribution.
134#[derive(Clone)]
135pub struct UserIntelligence {
136    service: Intelligence,
137    user_id: String,
138}
139
140enum AgentTurnBackend {
141    Codex(kcode_codex_runtime_v2::AgentTurn),
142    Buffered(BufferedAgentTurn),
143}
144
145struct BufferedAgentTurn {
146    events: VecDeque<kcode_codex_runtime_v2::AgentEvent>,
147    pending_call_id: Option<String>,
148    completed: Option<kcode_codex_runtime_v2::CompletedTurn>,
149}
150
151pub struct AgentTurn {
152    inner: AgentTurnBackend,
153    operation: ActiveOperation,
154    user_id: String,
155    requested_model: String,
156    actual_model: String,
157    provider_request_id: Option<String>,
158    operation_id: Uuid,
159    parent_operation_id: Option<Uuid>,
160    receipts: ReceiptStore,
161    receipt: Option<UsageReceipt>,
162}
163
164#[derive(Clone, Debug, Eq, PartialEq)]
165pub struct SearchRequest {
166    pub question: String,
167    pub model: String,
168    pub operation_id: Uuid,
169    pub parent_operation_id: Option<Uuid>,
170}
171
172#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
173#[serde(rename_all = "camelCase")]
174pub struct WebSource {
175    pub title: String,
176    pub url: String,
177}
178
179#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
180#[serde(rename_all = "camelCase")]
181pub struct SearchResponse {
182    pub answer: String,
183    pub sources: Vec<WebSource>,
184    pub model: String,
185    pub usage: Option<TokenUsage>,
186    pub cost: Option<CostEstimate>,
187}
188
189#[derive(Clone, Debug, Eq, PartialEq)]
190pub struct FetchRequest {
191    pub url: String,
192    pub operation_id: Uuid,
193    pub parent_operation_id: Option<Uuid>,
194}
195
196#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
197#[serde(rename_all = "camelCase")]
198pub struct FetchResponse {
199    pub url: String,
200    pub title: Option<String>,
201    pub content_type: String,
202    pub content: String,
203    pub truncated: bool,
204    pub retrieved_at: DateTime<Utc>,
205}
206
207#[derive(Clone, Copy, Debug, Eq, PartialEq)]
208pub enum MediaKind {
209    Image,
210    Audio,
211    Video,
212}
213
214#[derive(Clone, Debug, Eq, PartialEq)]
215pub struct Media {
216    pub kind: MediaKind,
217    pub bytes: Vec<u8>,
218    pub file_name: String,
219    pub content_type: String,
220}
221
222impl Media {
223    pub fn new(
224        kind: MediaKind,
225        bytes: Vec<u8>,
226        file_name: impl Into<String>,
227        content_type: impl Into<String>,
228    ) -> Result<Self> {
229        if bytes.is_empty() || bytes.len() > MAX_MEDIA_ANNOTATION_BYTES {
230            return Err(Error::invalid(format!(
231                "media must contain between 1 and {MAX_MEDIA_ANNOTATION_BYTES} bytes"
232            )));
233        }
234        let file_name = file_name.into();
235        let mut content_type = normalized_content_type(&content_type.into());
236        if kind == MediaKind::Audio && is_ogg(&file_name, &content_type) {
237            content_type = "audio/ogg".into();
238        }
239        Ok(Self {
240            kind,
241            bytes,
242            file_name,
243            content_type,
244        })
245    }
246
247    pub fn audio(
248        bytes: Vec<u8>,
249        file_name: impl Into<String>,
250        content_type: impl Into<String>,
251    ) -> Result<Self> {
252        Self::new(MediaKind::Audio, bytes, file_name, content_type)
253    }
254}
255
256#[derive(Clone, Debug, PartialEq)]
257pub struct TranscriptionRequest {
258    pub prompt: String,
259    pub model: String,
260    pub media: Media,
261    /// Optional sampling temperature. Omission preserves the provider default.
262    pub temperature: Option<f32>,
263    pub operation_id: Uuid,
264    pub parent_operation_id: Option<Uuid>,
265}
266
267#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
268#[serde(rename_all = "camelCase")]
269pub struct TranscriptionResponse {
270    pub model: String,
271    pub text: String,
272    pub metering: Metering,
273    pub cost: Option<CostEstimate>,
274}
275
276/// One Gemini audio-analysis call with exact model and accounting attribution.
277#[derive(Clone, Debug, PartialEq)]
278pub struct AudioAnalysisRequest {
279    pub operation: String,
280    pub prompt: String,
281    pub model: String,
282    pub media: Media,
283    pub schema: Option<Value>,
284    pub max_output_tokens: u32,
285    /// Optional sampling temperature. Omission preserves the provider default.
286    pub temperature: Option<f32>,
287    pub operation_id: Uuid,
288    pub parent_operation_id: Option<Uuid>,
289}
290
291#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
292#[serde(rename_all = "camelCase")]
293pub struct AudioAnalysisResponse {
294    pub model: String,
295    pub text: String,
296    pub usage: TokenUsage,
297    pub cost: CostEstimate,
298}
299
300/// Compatibility request for callers that always require structured Gemini audio output.
301#[derive(Clone, Debug, PartialEq)]
302pub struct StructuredAudioRequest {
303    pub operation: String,
304    pub prompt: String,
305    pub model: String,
306    pub media: Media,
307    pub schema: Value,
308    pub max_output_tokens: u32,
309    /// Optional sampling temperature. Omission preserves the provider default.
310    pub temperature: Option<f32>,
311    pub operation_id: Uuid,
312    pub parent_operation_id: Option<Uuid>,
313}
314
315#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
316#[serde(rename_all = "camelCase")]
317pub struct StructuredAudioResponse {
318    pub model: String,
319    pub text: String,
320    pub usage: TokenUsage,
321    pub cost: CostEstimate,
322}
323
324/// One tool-free Codex text generation call.
325#[derive(Clone, Debug, Eq, PartialEq)]
326pub struct TextGenerationRequest {
327    pub operation: String,
328    pub prompt: String,
329    pub model: String,
330    pub reasoning_effort: ReasoningEffort,
331    pub timeout: Duration,
332    pub operation_id: Uuid,
333    pub parent_operation_id: Option<Uuid>,
334}
335
336#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
337#[serde(rename_all = "camelCase")]
338pub struct TextGenerationResponse {
339    pub model: String,
340    pub text: String,
341    pub thread_id: String,
342    pub usage: Option<TokenUsage>,
343    pub cost: Option<CostEstimate>,
344}
345
346#[derive(Clone, Debug, Eq, PartialEq)]
347pub struct AnnotationRequest {
348    pub prompt: String,
349    pub model: String,
350    pub media: Media,
351    pub operation_id: Uuid,
352    pub parent_operation_id: Option<Uuid>,
353}
354
355#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
356#[serde(rename_all = "camelCase")]
357pub struct AnnotationResponse {
358    pub complete: bool,
359    pub model: String,
360    pub file_name: String,
361    pub content_type: String,
362    pub text: String,
363    pub incomplete_reason: Option<String>,
364    pub usage: Option<TokenUsage>,
365    pub cost: Option<CostEstimate>,
366}
367
368/// One direct image creation or reference-image modification request.
369#[derive(Clone, Debug, Eq, PartialEq)]
370pub struct ImageRequest {
371    /// Exact supported image model.
372    pub model: String,
373    /// Complete image generation or editing prompt.
374    pub prompt: String,
375    /// Existing images used as edit sources or visual references.
376    pub references: Vec<Media>,
377    /// Operation identifier used for cancellation.
378    pub operation_id: Uuid,
379    /// Optional running parent operation.
380    pub parent_operation_id: Option<Uuid>,
381}
382
383/// One generated image normalized across providers.
384#[derive(Clone, Debug, Eq, PartialEq)]
385pub struct ImageResponse {
386    /// Actual provider model.
387    pub model: String,
388    /// Output MIME type.
389    pub content_type: String,
390    /// Complete generated image bytes.
391    pub bytes: Vec<u8>,
392    /// Provider token usage, when returned.
393    pub usage: Option<TokenUsage>,
394    pub cost: Option<CostEstimate>,
395}
396
397#[derive(Clone, Debug, Eq, PartialEq)]
398pub struct Document {
399    pub bytes: Vec<u8>,
400    pub file_name: String,
401    pub content_type: String,
402}
403
404#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
405#[serde(rename_all = "camelCase")]
406pub struct DocumentExtraction {
407    pub file_name: String,
408    pub content_type: String,
409    pub format: String,
410    pub text: String,
411    pub characters: usize,
412    pub truncated: bool,
413}
414
415pub async fn open(config: Config) -> anyhow::Result<(Intelligence, RuntimeModel)> {
416    let mut codex_config = CodexConfig::new(DEFAULT_MODEL);
417    codex_config.base_instruction = KENNEDY_CODEX_BASE_INSTRUCTION.into();
418    codex_config.validation_reasoning_effort = GENERATION_REASONING_EFFORT;
419    let codex = Codex::open(codex_config, config.codex_catalog_cache)
420        .await
421        .context("opening Kennedy Codex runtime")?;
422    let limits = codex
423        .catalog()
424        .model_limits(DEFAULT_MODEL)
425        .with_context(|| format!("Codex model {DEFAULT_MODEL} is absent from the catalog"))?;
426    let runtime = RuntimeModel {
427        model: DEFAULT_MODEL.into(),
428        reasoning_effort: GENERATION_REASONING_EFFORT.as_str().into(),
429        context_window_tokens: limits.context_window_tokens(),
430        max_input_tokens: limits.max_input_tokens(),
431    };
432    let agent_config = kcode_codex_runtime_v2::CodexConfig {
433        executable: codex.catalog().executable().to_owned(),
434        base_instruction: KENNEDY_CODEX_BASE_INSTRUCTION.into(),
435        model_catalog: Some(codex.catalog().path().to_owned()),
436        ..kcode_codex_runtime_v2::CodexConfig::default()
437    };
438    let agent = kcode_codex_runtime_v2::Codex::open(agent_config)
439        .await
440        .context("opening Kennedy Codex runtime v2")?;
441    let openai = config
442        .openai_api_key
443        .filter(|value| !value.trim().is_empty())
444        .map(OpenAi::open)
445        .transpose()
446        .context("opening OpenAI client")?;
447    let gemini = config
448        .gemini_api_key
449        .filter(|value| !value.trim().is_empty())
450        .map(Gemini::open)
451        .transpose()
452        .context("opening Gemini client")?;
453    Ok((
454        Intelligence {
455            codex,
456            agent,
457            openai,
458            gemini,
459            web_fetcher: WebFetcher::default(),
460            document_extractor: DocumentExtractor::default(),
461            active_operations: OperationRegistry::default(),
462            receipts: ReceiptStore::open(config.receipt_directory)
463                .map_err(anyhow::Error::new)
464                .context("opening intelligence usage receipts")?,
465            model_cache: Arc::new(Mutex::new(HashMap::new())),
466        },
467        runtime,
468    ))
469}
470
471impl Intelligence {
472    pub fn for_user(&self, user_id: impl Into<String>) -> Result<UserIntelligence> {
473        let user_id = user_id.into();
474        if user_id.trim().is_empty() || user_id.chars().count() > 256 {
475            return Err(Error::invalid(
476                "user_id must contain between 1 and 256 characters",
477            ));
478        }
479        Ok(UserIntelligence {
480            service: self.clone(),
481            user_id,
482        })
483    }
484
485    pub fn cancel(&self, operation_id: Uuid) -> Result<bool> {
486        self.active_operations
487            .cancel(operation_id)
488            .map_err(registry_error)
489    }
490
491    pub fn receipts(&self) -> Result<Vec<UsageReceipt>> {
492        self.receipts.receipts()
493    }
494
495    pub fn daily_usage(
496        &self,
497        day: NaiveDate,
498    ) -> Result<std::collections::BTreeMap<DailyUsageKey, DailyUsage>> {
499        self.receipts.daily_usage(day)
500    }
501
502    pub async fn extract_document(&self, document: Document) -> Result<DocumentExtraction> {
503        extract_document(&self.document_extractor, document).await
504    }
505
506    /// Resolves an exact agent model through the configured provider boundary.
507    pub async fn resolve_agent_model(&self, requested: &str) -> Result<ResolvedAgentModel> {
508        validate_agent_model(requested)?;
509        if let Some(resolved) = self
510            .model_cache
511            .lock()
512            .map_err(|_| {
513                Error::internal("model_cache_unavailable", "The model cache is unavailable.")
514            })?
515            .get(requested)
516            .cloned()
517        {
518            return Ok(resolved);
519        }
520        let resolved = if let Some(provider_model) = requested.strip_prefix("codex/") {
521            let limits = self
522                .codex
523                .catalog()
524                .model_limits(provider_model)
525                .ok_or_else(|| {
526                    Error::invalid(format!(
527                        "{provider_model:?} is not an available model in the Codex catalog"
528                    ))
529                })?;
530            ResolvedAgentModel {
531                requested_model: requested.to_owned(),
532                provider_model: provider_model.to_owned(),
533                provider: AgentProvider::Codex,
534                context_window_tokens: limits.context_window_tokens(),
535                max_input_tokens: limits.max_input_tokens(),
536            }
537        } else if requested.ends_with("-sol")
538            || requested.ends_with("-terra")
539            || requested.ends_with("-luna")
540        {
541            let limits = self
542                .codex
543                .catalog()
544                .model_limits(requested)
545                .ok_or_else(|| {
546                    Error::invalid(format!(
547                        "{requested:?} is not an available model in the Codex catalog"
548                    ))
549                })?;
550            ResolvedAgentModel {
551                requested_model: requested.to_owned(),
552                provider_model: requested.to_owned(),
553                provider: AgentProvider::Codex,
554                context_window_tokens: limits.context_window_tokens(),
555                max_input_tokens: limits.max_input_tokens(),
556            }
557        } else if requested.starts_with("gemini-") {
558            let gemini = self.gemini.as_ref().ok_or_else(|| {
559                Error::unavailable("provider_not_configured", "Gemini is not configured.")
560            })?;
561            let metadata =
562                tokio::time::timeout(MODEL_DISCOVERY_TIMEOUT, gemini.model_metadata(requested))
563                    .await
564                    .map_err(|_| {
565                        Error::provider("provider_timeout", "Gemini model discovery timed out.")
566                    })?
567                    .map_err(gemini_error)?;
568            resolved_api_model(
569                requested,
570                metadata.id,
571                AgentProvider::Gemini,
572                metadata.context_window_tokens,
573                metadata.max_input_tokens,
574            )
575        } else {
576            let openai = self.openai.as_ref().ok_or_else(|| {
577                Error::unavailable("provider_not_configured", "OpenAI is not configured.")
578            })?;
579            let metadata =
580                tokio::time::timeout(MODEL_DISCOVERY_TIMEOUT, openai.model_metadata(requested))
581                    .await
582                    .map_err(|_| {
583                        Error::provider("provider_timeout", "OpenAI model discovery timed out.")
584                    })?
585                    .map_err(openai_error)?;
586            resolved_api_model(
587                requested,
588                metadata.id,
589                AgentProvider::OpenAi,
590                metadata.context_window_tokens,
591                metadata.max_input_tokens,
592            )
593        };
594        self.model_cache
595            .lock()
596            .map_err(|_| {
597                Error::internal("model_cache_unavailable", "The model cache is unavailable.")
598            })?
599            .insert(requested.to_owned(), resolved.clone());
600        Ok(resolved)
601    }
602}
603
604impl UserIntelligence {
605    pub fn user_id(&self) -> &str {
606        &self.user_id
607    }
608
609    pub async fn start_agent_turn(
610        &self,
611        operation_id: Uuid,
612        parent_operation_id: Option<Uuid>,
613        mut request: kcode_codex_runtime_v2::AgentRequest,
614    ) -> Result<AgentTurn> {
615        let resolved = self.service.resolve_agent_model(&request.model).await?;
616        if request.previous_thread_id.is_some() {
617            return Err(Error::invalid(
618                "agent turns must submit complete input without a previous thread; this keeps each usage receipt scoped to one call",
619            ));
620        }
621        let mut operation = self
622            .service
623            .active_operations
624            .register_request(operation_id, parent_operation_id)
625            .map_err(registry_error)?;
626        let requested_model = resolved.requested_model.clone();
627        request.model = resolved.provider_model.clone();
628        let (inner, actual_model, provider_request_id) = match resolved.provider {
629            AgentProvider::Codex => {
630                let inner = self.account_result(
631                    operation_id,
632                    parent_operation_id,
633                    "agent_turn",
634                    &requested_model,
635                    self.service
636                        .agent
637                        .start_turn(request)
638                        .await
639                        .map_err(codex_v2_error),
640                )?;
641                (
642                    AgentTurnBackend::Codex(inner),
643                    resolved.provider_model,
644                    None,
645                )
646            }
647            AgentProvider::OpenAi => {
648                let openai = self.service.openai.as_ref().ok_or_else(|| {
649                    Error::unavailable("provider_not_configured", "OpenAI is not configured.")
650                })?;
651                let provider_request = openai_agent_request(&request);
652                let provider_input = provider_request_json(&request);
653                let result = tokio::select! {
654                    _ = operation.cancelled() => Err(Error::cancelled()),
655                    result = tokio::time::timeout(request.timeout, openai.agent_turn(provider_request)) => {
656                        result
657                            .map_err(|_| Error::provider("provider_timeout", "OpenAI agent turn timed out."))
658                            .and_then(|result| result.map_err(openai_error))
659                    }
660                };
661                let result = self.account_result(
662                    operation_id,
663                    parent_operation_id,
664                    "agent_turn",
665                    &requested_model,
666                    result,
667                )?;
668                let actual_model = result.model.clone();
669                let provider_request_id = Some(result.response_id.clone());
670                (
671                    AgentTurnBackend::Buffered(buffered_openai_turn(provider_input, result)),
672                    actual_model,
673                    provider_request_id,
674                )
675            }
676            AgentProvider::Gemini => {
677                let gemini = self.service.gemini.as_ref().ok_or_else(|| {
678                    Error::unavailable("provider_not_configured", "Gemini is not configured.")
679                })?;
680                let provider_request = gemini_agent_request(&request);
681                let provider_input = provider_request_json(&request);
682                let result = tokio::select! {
683                    _ = operation.cancelled() => Err(Error::cancelled()),
684                    result = tokio::time::timeout(request.timeout, gemini.agent_turn(provider_request)) => {
685                        result
686                            .map_err(|_| Error::provider("provider_timeout", "Gemini agent turn timed out."))
687                            .and_then(|result| result.map_err(gemini_error))
688                    }
689                };
690                let result = self.account_result(
691                    operation_id,
692                    parent_operation_id,
693                    "agent_turn",
694                    &requested_model,
695                    result,
696                )?;
697                let actual_model = result.model.clone();
698                let provider_request_id = Some(result.interaction_id.clone());
699                (
700                    AgentTurnBackend::Buffered(buffered_gemini_turn(provider_input, result)),
701                    actual_model,
702                    provider_request_id,
703                )
704            }
705        };
706        Ok(AgentTurn {
707            inner,
708            operation,
709            user_id: self.user_id.clone(),
710            requested_model,
711            actual_model,
712            provider_request_id,
713            operation_id,
714            parent_operation_id,
715            receipts: self.service.receipts.clone(),
716            receipt: None,
717        })
718    }
719
720    pub async fn search(&self, request: SearchRequest) -> Result<Accounted<SearchResponse>> {
721        let question = request.question.trim();
722        if question.is_empty() || question.chars().count() > 4_000 {
723            return Err(Error::invalid(
724                "question must contain between 1 and 4000 characters",
725            ));
726        }
727        validate_model(&request.model)?;
728        let mut operation = self
729            .service
730            .active_operations
731            .register_request(request.operation_id, request.parent_operation_id)
732            .map_err(registry_error)?;
733        let started = Instant::now();
734        let (response, receipt) = if let Some(model) = gemini_model(&request.model) {
735            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
736                Error::unavailable(
737                    "provider_not_configured",
738                    "Gemini search is not configured.",
739                )
740            })?;
741            let result = tokio::select! {
742                _ = operation.cancelled() => Err(Error::cancelled()),
743                result = tokio::time::timeout(
744                    FAST_SEARCH_TIMEOUT,
745                    gemini.grounded_search_with_model(
746                        model,
747                        GroundedSearchRequest::new(question),
748                    ),
749                ) => result
750                    .map_err(|_| Error::provider("provider_timeout", "Gemini search timed out."))
751                    .and_then(|result| result.map_err(gemini_error)),
752            };
753            let result = self.account_result(
754                request.operation_id,
755                request.parent_operation_id,
756                "web_search",
757                &request.model,
758                result,
759            )?;
760            let interaction = result.interaction;
761            let usage = gemini_usage(&interaction.usage);
762            let cost = gemini_cost(&interaction.cost);
763            let receipt = self.record_tokens_with_cost(
764                request.operation_id,
765                request.parent_operation_id,
766                "web_search",
767                &request.model,
768                &interaction.model,
769                Some(usage),
770                Some(interaction.id.clone()),
771                None,
772                Some(cost.clone()),
773            )?;
774            let answer = interaction
775                .text
776                .filter(|value| !value.trim().is_empty())
777                .ok_or_else(|| {
778                    Error::provider("provider_error", "Gemini search returned no answer text.")
779                })
780                .map_err(|error| error.with_receipt(receipt.clone()))?;
781            (
782                SearchResponse {
783                    answer,
784                    sources: result
785                        .sources
786                        .into_iter()
787                        .map(|source| WebSource {
788                            title: source.title,
789                            url: source.url,
790                        })
791                        .collect(),
792                    model: interaction.model,
793                    usage: Some(usage),
794                    cost: Some(cost),
795                },
796                receipt,
797            )
798        } else {
799            let (reasoning, context, depth, timeout) = codex_search_profile(&request.model)?;
800            let result = tokio::select! {
801                _ = operation.cancelled() => Err(Error::cancelled()),
802                result = self.service.codex.web_search(CodexSearchRequest {
803                    question: question.to_owned(),
804                    model: request.model.clone(),
805                    reasoning_effort: reasoning,
806                    context,
807                    depth,
808                    timeout,
809                }) => result.map_err(codex_error),
810            };
811            let result = self.account_result(
812                request.operation_id,
813                request.parent_operation_id,
814                "web_search",
815                &request.model,
816                result,
817            )?;
818            let usage = result.usage.as_ref().map(codex_usage);
819            let cost = codex_search_cost(&request.model, usage);
820            let receipt = self.record_tokens_with_cost(
821                request.operation_id,
822                request.parent_operation_id,
823                "web_search",
824                &request.model,
825                &request.model,
826                usage,
827                None,
828                None,
829                cost.clone(),
830            )?;
831            (
832                SearchResponse {
833                    answer: result.answer,
834                    sources: result
835                        .sources
836                        .into_iter()
837                        .map(|source| WebSource {
838                            title: source.title,
839                            url: source.url,
840                        })
841                        .collect(),
842                    model: request.model.clone(),
843                    usage,
844                    cost,
845                },
846                receipt,
847            )
848        };
849        tracing::info!(
850            user_id = %self.user_id,
851            operation_id = %request.operation_id,
852            model = %response.model,
853            duration_ms = started.elapsed().as_millis(),
854            "Intelligence search completed"
855        );
856        Ok(Accounted {
857            value: response,
858            receipt,
859        })
860    }
861
862    /// Performs one Gemini audio-analysis call and records its provider usage.
863    pub async fn analyze_audio(
864        &self,
865        request: AudioAnalysisRequest,
866    ) -> Result<Accounted<AudioAnalysisResponse>> {
867        validate_operation(&request.operation)?;
868        validate_prompt(&request.prompt)?;
869        validate_model(&request.model)?;
870        if request.media.kind != MediaKind::Audio {
871            return Err(Error::invalid("audio analysis requires audio media"));
872        }
873        if request.max_output_tokens == 0 || request.max_output_tokens > 65_536 {
874            return Err(Error::invalid(
875                "max_output_tokens must be between 1 and 65536",
876            ));
877        }
878        let options = audio_generation_options(request.max_output_tokens, request.temperature)?;
879        let model = gemini_model(&request.model).ok_or_else(|| {
880            Error::invalid("audio analysis requires a supported exact Gemini model")
881        })?;
882        let gemini = self.service.gemini.as_ref().ok_or_else(|| {
883            Error::unavailable(
884                "provider_not_configured",
885                "Gemini audio analysis is not configured.",
886            )
887        })?;
888        let media = gemini_media(&request.media)?;
889        let mut provider_request = MultimodalRequest::new(request.prompt, vec![media]);
890        provider_request.options = options;
891        provider_request.structured_output = request
892            .schema
893            .map(StructuredOutput::new)
894            .transpose()
895            .map_err(gemini_error)?;
896        let mut operation = self
897            .service
898            .active_operations
899            .register_request(request.operation_id, request.parent_operation_id)
900            .map_err(registry_error)?;
901        let result = tokio::select! {
902            _ = operation.cancelled() => Err(Error::cancelled()),
903            result = tokio::time::timeout(
904                MEDIA_ANNOTATION_TIMEOUT,
905                gemini.infer_multimodal(model, provider_request),
906            ) => result
907                .map_err(|_| Error::provider("provider_timeout", "Gemini audio analysis timed out."))
908                .and_then(|result| result.map_err(gemini_error)),
909        };
910        let result = self.account_result(
911            request.operation_id,
912            request.parent_operation_id,
913            &request.operation,
914            &request.model,
915            result,
916        )?;
917        let usage = gemini_usage(&result.usage);
918        let cost = gemini_cost(&result.cost);
919        let receipt = self.record_tokens_with_cost(
920            request.operation_id,
921            request.parent_operation_id,
922            &request.operation,
923            &request.model,
924            &result.model,
925            Some(usage),
926            Some(result.id.clone()),
927            None,
928            Some(cost.clone()),
929        )?;
930        if result.status != GeminiCompletionStatus::Completed {
931            return Err(Error::provider(
932                "provider_incomplete",
933                "Gemini audio analysis did not complete.",
934            )
935            .with_receipt(receipt));
936        }
937        let text = result
938            .text
939            .filter(|text| !text.trim().is_empty())
940            .ok_or_else(|| {
941                Error::provider(
942                    "provider_empty_output",
943                    "Gemini audio analysis returned no text.",
944                )
945            })
946            .map_err(|error| error.with_receipt(receipt.clone()))?;
947        Ok(Accounted {
948            value: AudioAnalysisResponse {
949                model: result.model,
950                text,
951                usage,
952                cost,
953            },
954            receipt,
955        })
956    }
957
958    /// Performs one structured Gemini audio call and records its provider usage.
959    pub async fn transcribe_structured_audio(
960        &self,
961        request: StructuredAudioRequest,
962    ) -> Result<Accounted<StructuredAudioResponse>> {
963        self.analyze_audio(AudioAnalysisRequest {
964            operation: request.operation,
965            prompt: request.prompt,
966            model: request.model,
967            media: request.media,
968            schema: Some(request.schema),
969            max_output_tokens: request.max_output_tokens,
970            temperature: request.temperature,
971            operation_id: request.operation_id,
972            parent_operation_id: request.parent_operation_id,
973        })
974        .await
975        .map(|response| {
976            response.map(|response| StructuredAudioResponse {
977                model: response.model,
978                text: response.text,
979                usage: response.usage,
980                cost: response.cost,
981            })
982        })
983    }
984
985    /// Performs one tool-free Codex generation and records its provider usage.
986    pub async fn generate_text(
987        &self,
988        request: TextGenerationRequest,
989    ) -> Result<Accounted<TextGenerationResponse>> {
990        validate_operation(&request.operation)?;
991        validate_model(&request.model)?;
992        if request.prompt.trim().is_empty() || request.prompt.chars().count() > 1_000_000 {
993            return Err(Error::invalid(
994                "generation prompt must contain between 1 and 1000000 characters",
995            ));
996        }
997        if request.timeout.is_zero() || request.timeout > Duration::from_secs(7 * 60 * 60 + 30 * 60)
998        {
999            return Err(Error::invalid(
1000                "generation timeout must be between 1 second and 7 hours 30 minutes",
1001            ));
1002        }
1003        let mut operation = self
1004            .service
1005            .active_operations
1006            .register_request(request.operation_id, request.parent_operation_id)
1007            .map_err(registry_error)?;
1008        let mut provider_request =
1009            CodexGenerationRequest::new(request.prompt, request.model.clone());
1010        provider_request.reasoning_effort = request.reasoning_effort;
1011        provider_request.ephemeral = true;
1012        provider_request.timeout = request.timeout;
1013        let result = tokio::select! {
1014            _ = operation.cancelled() => Err(Error::cancelled()),
1015            result = self.service.codex.generate(provider_request) => result.map_err(codex_error),
1016        };
1017        let result = self.account_result(
1018            request.operation_id,
1019            request.parent_operation_id,
1020            &request.operation,
1021            &request.model,
1022            result,
1023        )?;
1024        let usage = result.usage.as_ref().map(codex_usage);
1025        let cost = usage.and_then(|usage| estimate_token_cost(&request.model, usage));
1026        let receipt = self.record_tokens(
1027            request.operation_id,
1028            request.parent_operation_id,
1029            &request.operation,
1030            &request.model,
1031            &request.model,
1032            usage,
1033            None,
1034            Some(result.thread_id.clone()),
1035        )?;
1036        Ok(Accounted {
1037            value: TextGenerationResponse {
1038                model: request.model,
1039                text: result.answer,
1040                thread_id: result.thread_id,
1041                usage,
1042                cost,
1043            },
1044            receipt,
1045        })
1046    }
1047
1048    pub async fn fetch(&self, request: FetchRequest) -> Result<FetchResponse> {
1049        let mut operation = self
1050            .service
1051            .active_operations
1052            .register_request(request.operation_id, request.parent_operation_id)
1053            .map_err(registry_error)?;
1054        let fetched = tokio::select! {
1055            _ = operation.cancelled() => Err(Error::cancelled()),
1056            result = self.service.web_fetcher.fetch(&request.url) => result.map_err(web_fetch_error),
1057        }?;
1058        Ok(FetchResponse {
1059            url: fetched.url,
1060            title: fetched.title,
1061            content_type: fetched.content_type,
1062            content: fetched.content,
1063            truncated: fetched.truncated,
1064            retrieved_at: DateTime::<Utc>::from(fetched.retrieved_at),
1065        })
1066    }
1067
1068    pub async fn transcribe(
1069        &self,
1070        request: TranscriptionRequest,
1071    ) -> Result<Accounted<TranscriptionResponse>> {
1072        validate_prompt(&request.prompt)?;
1073        if request.media.kind != MediaKind::Audio {
1074            return Err(Error::invalid("transcription requires audio media"));
1075        }
1076        validate_model(&request.model)?;
1077        let temperature = validate_transcription_temperature(&request.model, request.temperature)?;
1078        let mut operation = self
1079            .service
1080            .active_operations
1081            .register_request(request.operation_id, request.parent_operation_id)
1082            .map_err(registry_error)?;
1083        let (response, receipt) = if request.model == kcode_openai_api::GPT_4O_TRANSCRIBE {
1084            let openai = self.service.openai.as_ref().ok_or_else(|| {
1085                Error::unavailable(
1086                    "transcription_unavailable",
1087                    "OpenAI audio transcription is not configured.",
1088                )
1089            })?;
1090            let input = AudioInput::new(
1091                safe_audio_filename(&request.media.file_name, &request.media.content_type),
1092                request.media.content_type.clone(),
1093                request.media.bytes,
1094            )
1095            .map_err(openai_error)?;
1096            let mut provider_request = OpenAiTranscriptionRequest::new(input);
1097            provider_request.prompt = Some(request.prompt);
1098            let result = tokio::select! {
1099                _ = operation.cancelled() => Err(Error::cancelled()),
1100                result = tokio::time::timeout(
1101                    MEDIA_ANNOTATION_TIMEOUT,
1102                    openai.transcribe(provider_request),
1103                ) => result
1104                    .map_err(|_| Error::provider("provider_timeout", "OpenAI transcription timed out."))
1105                    .and_then(|result| result.map_err(openai_error)),
1106            };
1107            let result = self.account_result(
1108                request.operation_id,
1109                request.parent_operation_id,
1110                "transcribe_audio",
1111                &request.model,
1112                result,
1113            )?;
1114            let metering = result
1115                .usage
1116                .map(transcription_usage)
1117                .unwrap_or(Metering::Unavailable);
1118            let cost = estimate_cost(&request.model, &metering);
1119            let receipt = self.record_metering(
1120                request.operation_id,
1121                request.parent_operation_id,
1122                "transcribe_audio",
1123                &request.model,
1124                &request.model,
1125                metering.clone(),
1126                None,
1127                None,
1128            )?;
1129            (
1130                TranscriptionResponse {
1131                    model: request.model.clone(),
1132                    text: result.text,
1133                    metering,
1134                    cost,
1135                },
1136                receipt,
1137            )
1138        } else {
1139            let model = gemini_model(&request.model).ok_or_else(|| {
1140                Error::invalid("transcription model must be gpt-4o-transcribe or a supported exact Gemini model")
1141            })?;
1142            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1143                Error::unavailable(
1144                    "provider_not_configured",
1145                    "Gemini transcription is not configured.",
1146                )
1147            })?;
1148            let media = gemini_media(&request.media)?;
1149            let mut provider_request = MultimodalRequest::new(request.prompt, vec![media]);
1150            provider_request.options.temperature = temperature;
1151            let result = tokio::select! {
1152                _ = operation.cancelled() => Err(Error::cancelled()),
1153                result = tokio::time::timeout(
1154                    MEDIA_ANNOTATION_TIMEOUT,
1155                    gemini.infer_multimodal(model, provider_request),
1156                ) => result
1157                    .map_err(|_| Error::provider("provider_timeout", "Gemini transcription timed out."))
1158                    .and_then(|result| result.map_err(gemini_error)),
1159            };
1160            let result = self.account_result(
1161                request.operation_id,
1162                request.parent_operation_id,
1163                "transcribe_audio",
1164                &request.model,
1165                result,
1166            )?;
1167            let metering = Metering::Tokens(gemini_usage(&result.usage));
1168            let cost = gemini_cost(&result.cost);
1169            let receipt = self.record_metering_with_cost(
1170                request.operation_id,
1171                request.parent_operation_id,
1172                "transcribe_audio",
1173                &request.model,
1174                &result.model,
1175                metering.clone(),
1176                Some(result.id.clone()),
1177                None,
1178                Some(cost.clone()),
1179            )?;
1180            let text = result
1181                .text
1182                .filter(|text| !text.trim().is_empty())
1183                .ok_or_else(|| {
1184                    Error::provider(
1185                        "empty_transcription",
1186                        "Gemini returned no transcription text.",
1187                    )
1188                })
1189                .map_err(|error| error.with_receipt(receipt.clone()))?;
1190            (
1191                TranscriptionResponse {
1192                    model: result.model,
1193                    text,
1194                    metering,
1195                    cost: Some(cost),
1196                },
1197                receipt,
1198            )
1199        };
1200        Ok(Accounted {
1201            value: response,
1202            receipt,
1203        })
1204    }
1205
1206    pub async fn annotate(
1207        &self,
1208        request: AnnotationRequest,
1209    ) -> Result<Accounted<AnnotationResponse>> {
1210        validate_prompt(&request.prompt)?;
1211        validate_model(&request.model)?;
1212        let mut operation = self
1213            .service
1214            .active_operations
1215            .register_request(request.operation_id, request.parent_operation_id)
1216            .map_err(registry_error)?;
1217        let file_name = request.media.file_name.clone();
1218        let content_type = request.media.content_type.clone();
1219        let (response, receipt) = if let Some(model) = gemini_model(&request.model) {
1220            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1221                Error::unavailable(
1222                    "provider_not_configured",
1223                    "Gemini media annotation is not configured.",
1224                )
1225            })?;
1226            let media = gemini_media(&request.media)?;
1227            let result = tokio::select! {
1228                _ = operation.cancelled() => Err(Error::cancelled()),
1229                result = tokio::time::timeout(
1230                    MEDIA_ANNOTATION_TIMEOUT,
1231                    gemini.infer_multimodal(
1232                        model,
1233                        MultimodalRequest::new(request.prompt, vec![media]),
1234                    ),
1235                ) => result
1236                    .map_err(|_| Error::provider("provider_timeout", "Gemini annotation timed out."))
1237                    .and_then(|result| result.map_err(gemini_error)),
1238            };
1239            let result = self.account_result(
1240                request.operation_id,
1241                request.parent_operation_id,
1242                "annotate_media",
1243                &request.model,
1244                result,
1245            )?;
1246            let usage = gemini_usage(&result.usage);
1247            let cost = gemini_cost(&result.cost);
1248            let receipt = self.record_tokens_with_cost(
1249                request.operation_id,
1250                request.parent_operation_id,
1251                "annotate_media",
1252                &request.model,
1253                &result.model,
1254                Some(usage),
1255                Some(result.id.clone()),
1256                None,
1257                Some(cost.clone()),
1258            )?;
1259            let text = result
1260                .text
1261                .filter(|text| !text.trim().is_empty())
1262                .ok_or_else(|| {
1263                    Error::provider("empty_annotation", "Gemini returned no annotation text.")
1264                })
1265                .map_err(|error| error.with_receipt(receipt.clone()))?;
1266            (
1267                AnnotationResponse {
1268                    complete: result.status == GeminiCompletionStatus::Completed,
1269                    model: result.model,
1270                    file_name,
1271                    content_type,
1272                    text,
1273                    incomplete_reason: None,
1274                    usage: Some(usage),
1275                    cost: Some(cost),
1276                },
1277                receipt,
1278            )
1279        } else if request.model == "gpt-5.6" {
1280            if request.media.kind != MediaKind::Image {
1281                return Err(Error::invalid(
1282                    "OpenAI media annotation accepts images only",
1283                ));
1284            }
1285            let openai = self.service.openai.as_ref().ok_or_else(|| {
1286                Error::unavailable(
1287                    "provider_not_configured",
1288                    "OpenAI media annotation is not configured.",
1289                )
1290            })?;
1291            let image = OpenAiImageInput::new(
1292                openai_image_media_type(&request.media.content_type)?,
1293                request.media.bytes,
1294            )
1295            .map_err(openai_error)?;
1296            let result = tokio::select! {
1297                _ = operation.cancelled() => Err(Error::cancelled()),
1298                result = tokio::time::timeout(
1299                    MEDIA_ANNOTATION_TIMEOUT,
1300                    openai.analyze_image(ImageAnalysisRequest::new(image, request.prompt)),
1301                ) => result
1302                    .map_err(|_| Error::provider("provider_timeout", "OpenAI annotation timed out."))
1303                    .and_then(|result| result.map_err(openai_error)),
1304            };
1305            let result = self.account_result(
1306                request.operation_id,
1307                request.parent_operation_id,
1308                "annotate_media",
1309                &request.model,
1310                result,
1311            )?;
1312            let usage = result.usage.as_ref().map(openai_image_usage);
1313            let cost = result.usage.as_ref().map(openai_image_analysis_cost);
1314            let receipt = self.record_tokens_with_cost(
1315                request.operation_id,
1316                request.parent_operation_id,
1317                "annotate_media",
1318                &request.model,
1319                &result.model,
1320                usage,
1321                None,
1322                None,
1323                cost.clone(),
1324            )?;
1325            let (complete, incomplete_reason) = match result.status {
1326                OpenAiImageStatus::Completed => (true, None),
1327                OpenAiImageStatus::Incomplete { reason } => (false, reason),
1328            };
1329            (
1330                AnnotationResponse {
1331                    complete,
1332                    model: result.model,
1333                    file_name,
1334                    content_type,
1335                    text: result.text,
1336                    incomplete_reason,
1337                    usage,
1338                    cost,
1339                },
1340                receipt,
1341            )
1342        } else if matches!(
1343            request.model.as_str(),
1344            "gpt-5.6-sol" | "gpt-5.6-terra" | "gpt-5.6-luna"
1345        ) {
1346            if request.media.kind != MediaKind::Image {
1347                return Err(Error::invalid("Codex media annotation accepts images only"));
1348            }
1349            let image = kcode_codex_runtime_v2::ImageInput::new(
1350                codex_image_media_type(&request.media.content_type)?,
1351                request.media.bytes,
1352            )
1353            .map_err(codex_v2_error)?;
1354            let turn_result = self
1355                .service
1356                .agent
1357                .start_image_turn(kcode_codex_runtime_v2::ImageTurnRequest::new(
1358                    request.prompt,
1359                    request.model.clone(),
1360                    vec![image],
1361                ))
1362                .await
1363                .map_err(codex_v2_error);
1364            let mut turn = self.account_result(
1365                request.operation_id,
1366                request.parent_operation_id,
1367                "annotate_media",
1368                &request.model,
1369                turn_result,
1370            )?;
1371            let completed = loop {
1372                let event = tokio::select! {
1373                    _ = operation.cancelled() => {
1374                        turn.cancel();
1375                        let receipt = self.record_metering(
1376                            request.operation_id,
1377                            request.parent_operation_id,
1378                            "annotate_media",
1379                            &request.model,
1380                            &request.model,
1381                            Metering::Unavailable,
1382                            None,
1383                            None,
1384                        )?;
1385                        return Err(Error::cancelled().with_receipt(receipt));
1386                    }
1387                    event = turn.next_event() => event,
1388                };
1389                match event {
1390                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::ProviderInput(_))) => {}
1391                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::UsageUpdated(_))) => {}
1392                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::ToolCall(_))) => {
1393                        turn.cancel();
1394                        let receipt = self.record_metering(
1395                            request.operation_id,
1396                            request.parent_operation_id,
1397                            "annotate_media",
1398                            &request.model,
1399                            &request.model,
1400                            Metering::Unavailable,
1401                            None,
1402                            None,
1403                        )?;
1404                        return Err(Error::provider(
1405                            "provider_error",
1406                            "A tool-free Codex image turn requested a tool.",
1407                        )
1408                        .with_receipt(receipt));
1409                    }
1410                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::Completed(completed))) => {
1411                        break completed;
1412                    }
1413                    Some(Err(error)) => {
1414                        let receipt = self.record_metering(
1415                            request.operation_id,
1416                            request.parent_operation_id,
1417                            "annotate_media",
1418                            &request.model,
1419                            &request.model,
1420                            Metering::Unavailable,
1421                            None,
1422                            None,
1423                        )?;
1424                        return Err(codex_v2_error(error).with_receipt(receipt));
1425                    }
1426                    None => {
1427                        let receipt = self.record_metering(
1428                            request.operation_id,
1429                            request.parent_operation_id,
1430                            "annotate_media",
1431                            &request.model,
1432                            &request.model,
1433                            Metering::Unavailable,
1434                            None,
1435                            None,
1436                        )?;
1437                        return Err(Error::provider(
1438                            "empty_annotation",
1439                            "Codex ended without annotation text.",
1440                        )
1441                        .with_receipt(receipt));
1442                    }
1443                }
1444            };
1445            let usage = completed.usage.as_ref().map(codex_v2_usage);
1446            let cost = usage.and_then(|usage| estimate_token_cost(&request.model, usage));
1447            let receipt = self.record_tokens(
1448                request.operation_id,
1449                request.parent_operation_id,
1450                "annotate_media",
1451                &request.model,
1452                &request.model,
1453                usage,
1454                Some(completed.turn_id.clone()),
1455                Some(completed.thread_id.clone()),
1456            )?;
1457            (
1458                AnnotationResponse {
1459                    complete: true,
1460                    model: request.model.clone(),
1461                    file_name,
1462                    content_type,
1463                    text: completed.answer,
1464                    incomplete_reason: None,
1465                    usage,
1466                    cost,
1467                },
1468                receipt,
1469            )
1470        } else {
1471            return Err(Error::invalid(format!(
1472                "unsupported exact annotation model {}",
1473                request.model
1474            )));
1475        };
1476        Ok(Accounted {
1477            value: response,
1478            receipt,
1479        })
1480    }
1481
1482    /// Creates a new image or modifies supplied reference images.
1483    pub async fn generate_image(&self, request: ImageRequest) -> Result<Accounted<ImageResponse>> {
1484        validate_model(&request.model)?;
1485        if request.prompt.trim().is_empty()
1486            || request.prompt.chars().count() > MAX_IMAGE_PROMPT_CHARACTERS
1487        {
1488            return Err(Error::invalid(format!(
1489                "image prompt must contain 1 through {MAX_IMAGE_PROMPT_CHARACTERS} characters"
1490            )));
1491        }
1492        if request
1493            .references
1494            .iter()
1495            .any(|media| media.kind != MediaKind::Image)
1496        {
1497            return Err(Error::invalid("image references must all be images"));
1498        }
1499        let mut operation = self
1500            .service
1501            .active_operations
1502            .register_request(request.operation_id, request.parent_operation_id)
1503            .map_err(registry_error)?;
1504        let operation_name = if request.references.is_empty() {
1505            "generate_image"
1506        } else {
1507            "edit_image"
1508        };
1509        if request.model == kcode_openai_api::GPT_IMAGE_2 {
1510            let openai = self.service.openai.as_ref().ok_or_else(|| {
1511                Error::unavailable(
1512                    "provider_not_configured",
1513                    "OpenAI image generation is not configured.",
1514                )
1515            })?;
1516            let result = if request.references.is_empty() {
1517                let provider_request = OpenAiImageRequest::new(request.prompt);
1518                tokio::select! {
1519                    _ = operation.cancelled() => Err(Error::cancelled()),
1520                    result = tokio::time::timeout(
1521                        IMAGE_OPERATION_TIMEOUT,
1522                        openai.generate_image(provider_request),
1523                    ) => result
1524                        .map_err(|_| Error::provider("provider_timeout", "OpenAI image generation timed out."))
1525                        .and_then(|result| result.map_err(openai_error)),
1526                }
1527            } else {
1528                let mut images = request
1529                    .references
1530                    .into_iter()
1531                    .map(|media| {
1532                        OpenAiImageInput::new(
1533                            openai_image_media_type(&media.content_type)?,
1534                            media.bytes,
1535                        )
1536                        .map_err(openai_error)
1537                    })
1538                    .collect::<Result<Vec<_>>>()?;
1539                let first = images.remove(0);
1540                let mut provider_request = OpenAiImageEditRequest::new(first, request.prompt);
1541                provider_request.images.extend(images);
1542                tokio::select! {
1543                    _ = operation.cancelled() => Err(Error::cancelled()),
1544                    result = tokio::time::timeout(
1545                        IMAGE_OPERATION_TIMEOUT,
1546                        openai.edit_image(provider_request),
1547                    ) => result
1548                        .map_err(|_| Error::provider("provider_timeout", "OpenAI image editing timed out."))
1549                        .and_then(|result| result.map_err(openai_error)),
1550                }
1551            };
1552            let result = self.account_result(
1553                request.operation_id,
1554                request.parent_operation_id,
1555                operation_name,
1556                &request.model,
1557                result,
1558            )?;
1559            let usage = result.usage.as_ref().map(openai_generation_usage);
1560            let cost = result.usage.as_ref().map(openai_generation_cost);
1561            let receipt = self.record_tokens_with_cost(
1562                request.operation_id,
1563                request.parent_operation_id,
1564                operation_name,
1565                &request.model,
1566                kcode_openai_api::GPT_IMAGE_2,
1567                usage,
1568                result.request_id,
1569                None,
1570                cost.clone(),
1571            )?;
1572            Ok(Accounted {
1573                value: ImageResponse {
1574                    model: kcode_openai_api::GPT_IMAGE_2.into(),
1575                    content_type: result.image.format.mime_type().into(),
1576                    bytes: result.image.data,
1577                    usage,
1578                    cost,
1579                },
1580                receipt,
1581            })
1582        } else if request.model == kcode_gemini_api::NANO_BANANA_PRO {
1583            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1584                Error::unavailable(
1585                    "provider_not_configured",
1586                    "Gemini image generation is not configured.",
1587                )
1588            })?;
1589            let mut provider_request = NanoBananaProRequest::new(request.prompt);
1590            provider_request.images = request
1591                .references
1592                .into_iter()
1593                .map(|media| {
1594                    GeminiMediaInput::image(&media.content_type, media.bytes).map_err(gemini_error)
1595                })
1596                .collect::<Result<Vec<_>>>()?;
1597            let result = tokio::select! {
1598                _ = operation.cancelled() => Err(Error::cancelled()),
1599                result = tokio::time::timeout(
1600                    IMAGE_OPERATION_TIMEOUT,
1601                    gemini.nano_banana_pro(provider_request),
1602                ) => result
1603                    .map_err(|_| Error::provider("provider_timeout", "Gemini image generation timed out."))
1604                    .and_then(|result| result.map_err(gemini_error)),
1605            };
1606            let result = self.account_result(
1607                request.operation_id,
1608                request.parent_operation_id,
1609                operation_name,
1610                &request.model,
1611                result,
1612            )?;
1613            let usage = gemini_usage(&result.usage);
1614            let cost = gemini_cost(&result.cost);
1615            let receipt = self.record_tokens_with_cost(
1616                request.operation_id,
1617                request.parent_operation_id,
1618                operation_name,
1619                &request.model,
1620                &result.model,
1621                Some(usage),
1622                Some(result.id.clone()),
1623                None,
1624                Some(cost.clone()),
1625            )?;
1626            let mut images = result.images;
1627            if images.len() != 1 {
1628                return Err(Error::provider(
1629                    "provider_error",
1630                    "Gemini did not return exactly one generated image.",
1631                )
1632                .with_receipt(receipt));
1633            }
1634            let image = images.remove(0);
1635            Ok(Accounted {
1636                value: ImageResponse {
1637                    model: result.model,
1638                    content_type: image.mime_type,
1639                    bytes: image.data,
1640                    usage: Some(usage),
1641                    cost: Some(cost),
1642                },
1643                receipt,
1644            })
1645        } else {
1646            Err(Error::invalid(format!(
1647                "unsupported exact image model {}",
1648                request.model
1649            )))
1650        }
1651    }
1652
1653    #[allow(clippy::too_many_arguments)]
1654    fn record_tokens(
1655        &self,
1656        operation_id: Uuid,
1657        parent_operation_id: Option<Uuid>,
1658        operation: &str,
1659        requested_model: &str,
1660        actual_model: &str,
1661        usage: Option<TokenUsage>,
1662        provider_request_id: Option<String>,
1663        provider_thread_id: Option<String>,
1664    ) -> Result<UsageReceipt> {
1665        self.record_metering_with_cost(
1666            operation_id,
1667            parent_operation_id,
1668            operation,
1669            requested_model,
1670            actual_model,
1671            usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
1672            provider_request_id,
1673            provider_thread_id,
1674            None,
1675        )
1676    }
1677
1678    #[allow(clippy::too_many_arguments)]
1679    fn record_tokens_with_cost(
1680        &self,
1681        operation_id: Uuid,
1682        parent_operation_id: Option<Uuid>,
1683        operation: &str,
1684        requested_model: &str,
1685        actual_model: &str,
1686        usage: Option<TokenUsage>,
1687        provider_request_id: Option<String>,
1688        provider_thread_id: Option<String>,
1689        cost: Option<CostEstimate>,
1690    ) -> Result<UsageReceipt> {
1691        self.record_metering_with_cost(
1692            operation_id,
1693            parent_operation_id,
1694            operation,
1695            requested_model,
1696            actual_model,
1697            usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
1698            provider_request_id,
1699            provider_thread_id,
1700            cost,
1701        )
1702    }
1703
1704    #[allow(clippy::too_many_arguments)]
1705    fn record_metering(
1706        &self,
1707        operation_id: Uuid,
1708        parent_operation_id: Option<Uuid>,
1709        operation: &str,
1710        requested_model: &str,
1711        actual_model: &str,
1712        metering: Metering,
1713        provider_request_id: Option<String>,
1714        provider_thread_id: Option<String>,
1715    ) -> Result<UsageReceipt> {
1716        self.record_metering_with_cost(
1717            operation_id,
1718            parent_operation_id,
1719            operation,
1720            requested_model,
1721            actual_model,
1722            metering,
1723            provider_request_id,
1724            provider_thread_id,
1725            None,
1726        )
1727    }
1728
1729    #[allow(clippy::too_many_arguments)]
1730    fn record_metering_with_cost(
1731        &self,
1732        operation_id: Uuid,
1733        parent_operation_id: Option<Uuid>,
1734        operation: &str,
1735        requested_model: &str,
1736        actual_model: &str,
1737        metering: Metering,
1738        provider_request_id: Option<String>,
1739        provider_thread_id: Option<String>,
1740        cost: Option<CostEstimate>,
1741    ) -> Result<UsageReceipt> {
1742        let mut receipt = UsageReceipt::new(
1743            self.user_id.clone(),
1744            operation_id,
1745            parent_operation_id,
1746            operation,
1747            requested_model,
1748            actual_model,
1749            metering,
1750        );
1751        receipt.provider_request_id = provider_request_id;
1752        receipt.provider_thread_id = provider_thread_id;
1753        if cost.is_some() {
1754            receipt.cost = cost;
1755        }
1756        self.service.receipts.record(&receipt)?;
1757        Ok(receipt)
1758    }
1759
1760    fn account_result<T>(
1761        &self,
1762        operation_id: Uuid,
1763        parent_operation_id: Option<Uuid>,
1764        operation: &str,
1765        requested_model: &str,
1766        result: Result<T>,
1767    ) -> Result<T> {
1768        match result {
1769            Ok(value) => Ok(value),
1770            Err(error) => {
1771                let receipt = self.record_metering(
1772                    operation_id,
1773                    parent_operation_id,
1774                    operation,
1775                    requested_model,
1776                    requested_model,
1777                    Metering::Unavailable,
1778                    None,
1779                    None,
1780                )?;
1781                Err(error.with_receipt(receipt))
1782            }
1783        }
1784    }
1785}
1786
1787async fn extract_document(
1788    extractor: &DocumentExtractor,
1789    document: Document,
1790) -> Result<DocumentExtraction> {
1791    let extractor = extractor.clone();
1792    let extracted = tokio::task::spawn_blocking(move || {
1793        extractor.extract(DocumentInput {
1794            file_name: document.file_name,
1795            content_type: document.content_type,
1796            data: document.bytes,
1797        })
1798    })
1799    .await
1800    .map_err(|_| {
1801        Error::internal(
1802            "document_extraction_failed",
1803            "The document extraction worker stopped unexpectedly.",
1804        )
1805    })?
1806    .map_err(document_error)?;
1807    Ok(DocumentExtraction {
1808        file_name: extracted.file_name,
1809        content_type: extracted.content_type,
1810        format: extracted.format.as_str().into(),
1811        text: extracted.text,
1812        characters: extracted.characters,
1813        truncated: extracted.truncated,
1814    })
1815}
1816
1817impl AgentTurn {
1818    /// Canonical receipt after the turn has completed or been marked unavailable.
1819    pub fn receipt(&self) -> Option<&UsageReceipt> {
1820        self.receipt.as_ref()
1821    }
1822
1823    /// Ends an interrupted turn with unavailable metering and returns its receipt.
1824    pub fn finish_unavailable(&mut self) -> Result<&UsageReceipt> {
1825        if self.receipt.is_none() {
1826            self.record_unavailable()?;
1827        }
1828        Ok(self.receipt.as_ref().expect("receipt was just recorded"))
1829    }
1830
1831    pub async fn next_event(&mut self) -> Result<Option<kcode_codex_runtime_v2::AgentEvent>> {
1832        let event = match &mut self.inner {
1833            AgentTurnBackend::Codex(inner) => tokio::select! {
1834                _ = self.operation.cancelled() => {
1835                    inner.cancel();
1836                    let receipt = self.record_unavailable()?;
1837                    return Err(Error::cancelled().with_receipt(receipt));
1838                }
1839                event = inner.next_event() => match event {
1840                    Some(Ok(event)) => Some(event),
1841                    Some(Err(error)) => {
1842                        let receipt = self.record_unavailable()?;
1843                        return Err(codex_v2_error(error).with_receipt(receipt));
1844                    }
1845                    None => {
1846                        self.record_unavailable()?;
1847                        None
1848                    },
1849                }
1850            },
1851            AgentTurnBackend::Buffered(inner) => {
1852                if self.operation.direct_cancellation_requested() {
1853                    let receipt = self.record_unavailable()?;
1854                    return Err(Error::cancelled().with_receipt(receipt));
1855                }
1856                inner.events.pop_front()
1857            }
1858        };
1859        if let Some(kcode_codex_runtime_v2::AgentEvent::Completed(completed)) = &event
1860            && self.receipt.is_none()
1861        {
1862            let mut receipt = UsageReceipt::new(
1863                self.user_id.clone(),
1864                self.operation_id,
1865                self.parent_operation_id,
1866                "agent_turn",
1867                self.requested_model.clone(),
1868                self.actual_model.clone(),
1869                completed
1870                    .usage
1871                    .as_ref()
1872                    .map(codex_v2_usage)
1873                    .map(Metering::Tokens)
1874                    .unwrap_or(Metering::Unavailable),
1875            );
1876            receipt.provider_request_id = self
1877                .provider_request_id
1878                .clone()
1879                .or_else(|| Some(completed.turn_id.clone()));
1880            receipt.provider_thread_id = Some(completed.thread_id.clone());
1881            self.receipts.record(&receipt)?;
1882            self.receipt = Some(receipt);
1883        }
1884        Ok(event)
1885    }
1886
1887    fn record_unavailable(&mut self) -> Result<UsageReceipt> {
1888        if let Some(receipt) = &self.receipt {
1889            return Ok(receipt.clone());
1890        }
1891        let receipt = UsageReceipt::new(
1892            self.user_id.clone(),
1893            self.operation_id,
1894            self.parent_operation_id,
1895            "agent_turn",
1896            self.requested_model.clone(),
1897            self.actual_model.clone(),
1898            Metering::Unavailable,
1899        );
1900        self.receipts.record(&receipt)?;
1901        self.receipt = Some(receipt.clone());
1902        Ok(receipt)
1903    }
1904
1905    pub async fn respond(
1906        &mut self,
1907        call_id: &str,
1908        result: kcode_codex_runtime_v2::ToolResult,
1909    ) -> Result<()> {
1910        match &mut self.inner {
1911            AgentTurnBackend::Codex(inner) => {
1912                inner.respond(call_id, result).await.map_err(codex_v2_error)
1913            }
1914            AgentTurnBackend::Buffered(inner) => {
1915                if inner.pending_call_id.as_deref() != Some(call_id) {
1916                    return Err(Error::invalid(
1917                        "tool result does not match the pending provider call",
1918                    ));
1919                }
1920                inner.pending_call_id = None;
1921                if let Some(mut completed) = inner.completed.take() {
1922                    completed.answer.clear();
1923                    inner
1924                        .events
1925                        .push_back(kcode_codex_runtime_v2::AgentEvent::Completed(completed));
1926                }
1927                Ok(())
1928            }
1929        }
1930    }
1931}
1932
1933impl Drop for AgentTurn {
1934    fn drop(&mut self) {
1935        let _ = self.record_unavailable();
1936    }
1937}
1938
1939fn registry_error(error: RegistryError) -> Error {
1940    match error {
1941        RegistryError::RegistryUnavailable => Error::internal(
1942            "operation_registry_unavailable",
1943            "The operation registry is unavailable.",
1944        ),
1945        RegistryError::SameAsParent => {
1946            Error::invalid("operation_id and parent_operation_id must be different")
1947        }
1948        RegistryError::OperationInProgress => Error::conflict(
1949            "operation_in_progress",
1950            "An operation with this identifier is already running.",
1951        ),
1952        RegistryError::ParentNotRunning => Error::conflict(
1953            "parent_operation_not_running",
1954            "The parent operation is no longer running.",
1955        ),
1956    }
1957}
1958
1959fn validate_model(model: &str) -> Result<()> {
1960    if model.trim().is_empty()
1961        || model.chars().count() > 128
1962        || !model
1963            .bytes()
1964            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_'))
1965    {
1966        return Err(Error::invalid(
1967            "model must be an exact safe model identifier",
1968        ));
1969    }
1970    Ok(())
1971}
1972
1973fn validate_agent_model(model: &str) -> Result<()> {
1974    if model.trim().is_empty()
1975        || model.chars().count() > 128
1976        || !model.bytes().all(|byte| {
1977            byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_' | b'/' | b':')
1978        })
1979        || model.starts_with("codex:")
1980    {
1981        return Err(Error::invalid(
1982            "model must be an exact safe provider model identifier",
1983        ));
1984    }
1985    Ok(())
1986}
1987
1988fn validate_operation(operation: &str) -> Result<()> {
1989    if operation.is_empty()
1990        || operation.len() > 64
1991        || !operation
1992            .bytes()
1993            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
1994    {
1995        return Err(Error::invalid(
1996            "operation must be a lowercase identifier of at most 64 bytes",
1997        ));
1998    }
1999    Ok(())
2000}
2001
2002fn validate_prompt(prompt: &str) -> Result<()> {
2003    if prompt.trim().is_empty() {
2004        return Err(Error::invalid("prompt must not be blank"));
2005    }
2006    Ok(())
2007}
2008
2009fn validate_temperature(temperature: Option<f32>) -> Result<Option<f32>> {
2010    match temperature {
2011        Some(value) if value.is_finite() && (0.0..=2.0).contains(&value) => Ok(Some(value)),
2012        Some(_) => Err(Error::invalid(
2013            "temperature must be finite and between 0.0 and 2.0 inclusive",
2014        )),
2015        None => Ok(None),
2016    }
2017}
2018
2019fn validate_transcription_temperature(
2020    model: &str,
2021    temperature: Option<f32>,
2022) -> Result<Option<f32>> {
2023    let temperature = validate_temperature(temperature)?;
2024    if model == kcode_openai_api::GPT_4O_TRANSCRIBE && temperature.is_some() {
2025        return Err(Error::invalid(
2026            "OpenAI transcription does not accept a temperature",
2027        ));
2028    }
2029    Ok(temperature)
2030}
2031
2032fn audio_generation_options(
2033    max_output_tokens: u32,
2034    temperature: Option<f32>,
2035) -> Result<GenerationOptions> {
2036    Ok(GenerationOptions {
2037        max_output_tokens: Some(max_output_tokens),
2038        temperature: validate_temperature(temperature)?,
2039        thinking_level: Some(ThinkingLevel::High),
2040        service_tier: ServiceTier::Standard,
2041    })
2042}
2043
2044fn gemini_model(model: &str) -> Option<TextModel> {
2045    match model {
2046        kcode_gemini_api::GEMINI_25_FLASH => Some(TextModel::Flash25),
2047        kcode_gemini_api::GEMINI_31_FLASH_LITE => Some(TextModel::FlashLite),
2048        kcode_gemini_api::GEMINI_31_PRO => Some(TextModel::Pro),
2049        _ => None,
2050    }
2051}
2052
2053fn codex_search_profile(
2054    model: &str,
2055) -> Result<(
2056    kcode_codex_runtime::ReasoningEffort,
2057    kcode_codex_runtime::WebSearchContext,
2058    kcode_codex_runtime::SearchDepth,
2059    Duration,
2060)> {
2061    match model {
2062        QUALITY_SEARCH_MODEL => Ok((
2063            QUALITY_SEARCH_REASONING,
2064            QUALITY_SEARCH_CONTEXT,
2065            QUALITY_SEARCH_DEPTH,
2066            QUALITY_SEARCH_TIMEOUT,
2067        )),
2068        BALANCED_SEARCH_MODEL => Ok((
2069            BALANCED_SEARCH_REASONING,
2070            BALANCED_SEARCH_CONTEXT,
2071            BALANCED_SEARCH_DEPTH,
2072            BALANCED_SEARCH_TIMEOUT,
2073        )),
2074        _ => Err(Error::invalid(
2075            "unsupported exact web-search model; use a supported Gemini model, gpt-5.6-sol, or gpt-5.6-terra",
2076        )),
2077    }
2078}
2079
2080fn normalized_content_type(value: &str) -> String {
2081    value
2082        .split(';')
2083        .next()
2084        .unwrap_or("application/octet-stream")
2085        .trim()
2086        .to_ascii_lowercase()
2087}
2088
2089fn is_ogg(file_name: &str, content_type: &str) -> bool {
2090    matches!(content_type, "audio/ogg" | "video/ogg" | "application/ogg")
2091        || file_name.rsplit_once('.').is_some_and(|(_, extension)| {
2092            matches!(
2093                extension.to_ascii_lowercase().as_str(),
2094                "ogg" | "oga" | "opus"
2095            )
2096        })
2097}
2098
2099fn safe_audio_filename(value: &str, content_type: &str) -> String {
2100    let extension = match content_type {
2101        "audio/ogg" | "audio/opus" | "application/ogg" | "video/ogg" => "ogg",
2102        "audio/wav" | "audio/x-wav" => "wav",
2103        "audio/mpeg" | "audio/mp3" => "mp3",
2104        "audio/mp4" => "mp4",
2105        "audio/webm" => "webm",
2106        "audio/flac" | "audio/x-flac" => "flac",
2107        "audio/m4a" => "m4a",
2108        _ => "audio",
2109    };
2110    let cleaned = value
2111        .chars()
2112        .filter(|character| {
2113            character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_')
2114        })
2115        .take(120)
2116        .collect::<String>();
2117    let supported = cleaned.rsplit_once('.').is_some_and(|(_, extension)| {
2118        matches!(
2119            extension.to_ascii_lowercase().as_str(),
2120            "flac"
2121                | "mp3"
2122                | "mp4"
2123                | "mpeg"
2124                | "mpga"
2125                | "m4a"
2126                | "ogg"
2127                | "oga"
2128                | "opus"
2129                | "wav"
2130                | "webm"
2131        )
2132    });
2133    if cleaned.is_empty() || !supported {
2134        format!("voice-note.{extension}")
2135    } else {
2136        cleaned
2137    }
2138}
2139
2140fn gemini_media(media: &Media) -> Result<GeminiMediaInput> {
2141    match media.kind {
2142        MediaKind::Image => GeminiMediaInput::image(&media.content_type, media.bytes.clone()),
2143        MediaKind::Audio => GeminiMediaInput::audio(&media.content_type, media.bytes.clone()),
2144        MediaKind::Video => GeminiMediaInput::video(&media.content_type, media.bytes.clone()),
2145    }
2146    .map_err(gemini_error)
2147}
2148
2149fn resolved_api_model(
2150    requested: &str,
2151    provider_model: String,
2152    provider: AgentProvider,
2153    context_window_tokens: Option<u64>,
2154    max_input_tokens: Option<u64>,
2155) -> ResolvedAgentModel {
2156    let known = match provider {
2157        AgentProvider::OpenAi if requested == "gpt-5.6" => Some((1_000_000, 700_000)),
2158        AgentProvider::Gemini
2159            if matches!(
2160                requested,
2161                "gemini-2.5-flash" | "gemini-3.1-flash-lite" | "gemini-3.1-pro-preview"
2162            ) =>
2163        {
2164            Some((1_000_000, 700_000))
2165        }
2166        _ => None,
2167    };
2168    let context_window_tokens = context_window_tokens
2169        .or_else(|| known.map(|limits| limits.0))
2170        .unwrap_or(128_000);
2171    let max_input_tokens = max_input_tokens
2172        .or_else(|| known.map(|limits| limits.1))
2173        .unwrap_or_else(|| context_window_tokens.saturating_mul(70) / 100)
2174        .min(context_window_tokens);
2175    ResolvedAgentModel {
2176        requested_model: requested.to_owned(),
2177        provider_model,
2178        provider,
2179        context_window_tokens,
2180        max_input_tokens,
2181    }
2182}
2183
2184fn openai_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> OpenAiAgentRequest {
2185    OpenAiAgentRequest {
2186        model: request.model.clone(),
2187        input: request.input.clone(),
2188        reasoning_effort: request.reasoning_effort.as_str().into(),
2189        tools: request
2190            .tools
2191            .iter()
2192            .map(|tool| OpenAiAgentTool {
2193                name: tool.name.clone(),
2194                description: tool.description.clone(),
2195                input_schema: tool.input_schema.clone(),
2196            })
2197            .collect(),
2198    }
2199}
2200
2201fn gemini_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> GeminiAgentRequest {
2202    GeminiAgentRequest {
2203        model: request.model.clone(),
2204        input: request.input.clone(),
2205        reasoning_effort: request.reasoning_effort.as_str().into(),
2206        tools: request
2207            .tools
2208            .iter()
2209            .map(|tool| GeminiAgentTool {
2210                name: tool.name.clone(),
2211                description: tool.description.clone(),
2212                input_schema: tool.input_schema.clone(),
2213            })
2214            .collect(),
2215    }
2216}
2217
2218fn provider_request_json(request: &kcode_codex_runtime_v2::AgentRequest) -> String {
2219    format!(
2220        "{}\n",
2221        serde_json::to_string(&json!({
2222            "model": request.model,
2223            "input": request.input,
2224            "reasoningEffort": request.reasoning_effort.as_str(),
2225            "tools": request.tools.iter().map(|tool| json!({
2226                "name": tool.name,
2227                "description": tool.description,
2228                "parameters": tool.input_schema
2229            })).collect::<Vec<_>>()
2230        }))
2231        .expect("provider request values always serialize")
2232    )
2233}
2234
2235fn buffered_openai_turn(
2236    provider_input: String,
2237    response: kcode_openai_api::AgentTurnResponse,
2238) -> BufferedAgentTurn {
2239    let usage = response
2240        .usage
2241        .as_ref()
2242        .map(|usage| kcode_codex_runtime_v2::TokenUsage {
2243            input_tokens: usage.input_tokens,
2244            output_tokens: usage.output_tokens,
2245            cached_input_tokens: usage.cached_input_tokens,
2246            reasoning_output_tokens: usage.reasoning_output_tokens,
2247            last_input_tokens: None,
2248            last_output_tokens: None,
2249        });
2250    let completed = kcode_codex_runtime_v2::CompletedTurn {
2251        thread_id: response.response_id,
2252        turn_id: Uuid::new_v4().to_string(),
2253        answer: response.text,
2254        usage,
2255    };
2256    buffered_turn(
2257        provider_input,
2258        response
2259            .tool_call
2260            .map(|call| kcode_codex_runtime_v2::DynamicToolCall {
2261                call_id: call.call_id,
2262                tool: call.name,
2263                arguments: call.arguments,
2264            }),
2265        completed,
2266    )
2267}
2268
2269fn buffered_gemini_turn(
2270    provider_input: String,
2271    response: kcode_gemini_api::AgentTurnResponse,
2272) -> BufferedAgentTurn {
2273    let usage = kcode_codex_runtime_v2::TokenUsage {
2274        input_tokens: response.usage.input_tokens,
2275        output_tokens: response
2276            .usage
2277            .output_tokens
2278            .saturating_add(response.usage.thought_tokens),
2279        cached_input_tokens: response.usage.cached_tokens,
2280        reasoning_output_tokens: response.usage.thought_tokens,
2281        last_input_tokens: None,
2282        last_output_tokens: None,
2283    };
2284    let completed = kcode_codex_runtime_v2::CompletedTurn {
2285        thread_id: response.interaction_id,
2286        turn_id: Uuid::new_v4().to_string(),
2287        answer: response.text,
2288        usage: Some(usage),
2289    };
2290    buffered_turn(
2291        provider_input,
2292        response
2293            .tool_call
2294            .map(|call| kcode_codex_runtime_v2::DynamicToolCall {
2295                call_id: call.call_id,
2296                tool: call.name,
2297                arguments: call.arguments,
2298            }),
2299        completed,
2300    )
2301}
2302
2303fn buffered_turn(
2304    provider_input: String,
2305    tool_call: Option<kcode_codex_runtime_v2::DynamicToolCall>,
2306    completed: kcode_codex_runtime_v2::CompletedTurn,
2307) -> BufferedAgentTurn {
2308    let mut events = VecDeque::from([kcode_codex_runtime_v2::AgentEvent::ProviderInput(
2309        provider_input,
2310    )]);
2311    if let Some(usage) = completed.usage.clone() {
2312        events.push_back(kcode_codex_runtime_v2::AgentEvent::UsageUpdated(usage));
2313    }
2314    let pending_call_id = tool_call.as_ref().map(|call| call.call_id.clone());
2315    if let Some(call) = tool_call {
2316        events.push_back(kcode_codex_runtime_v2::AgentEvent::ToolCall(call));
2317    } else {
2318        events.push_back(kcode_codex_runtime_v2::AgentEvent::Completed(
2319            completed.clone(),
2320        ));
2321    }
2322    BufferedAgentTurn {
2323        events,
2324        pending_call_id,
2325        completed: Some(completed),
2326    }
2327}
2328
2329fn openai_image_media_type(content_type: &str) -> Result<OpenAiImageMediaType> {
2330    match content_type {
2331        "image/png" => Ok(OpenAiImageMediaType::Png),
2332        "image/jpeg" | "image/jpg" => Ok(OpenAiImageMediaType::Jpeg),
2333        "image/webp" => Ok(OpenAiImageMediaType::WebP),
2334        "image/gif" => Ok(OpenAiImageMediaType::Gif),
2335        _ => Err(Error::invalid(
2336            "OpenAI annotations require PNG, JPEG, WebP, or GIF",
2337        )),
2338    }
2339}
2340
2341fn codex_image_media_type(content_type: &str) -> Result<kcode_codex_runtime_v2::ImageMediaType> {
2342    match content_type {
2343        "image/png" => Ok(kcode_codex_runtime_v2::ImageMediaType::Png),
2344        "image/jpeg" | "image/jpg" => Ok(kcode_codex_runtime_v2::ImageMediaType::Jpeg),
2345        "image/webp" => Ok(kcode_codex_runtime_v2::ImageMediaType::Webp),
2346        _ => Err(Error::invalid(
2347            "Codex annotations require PNG, JPEG, or WebP",
2348        )),
2349    }
2350}
2351
2352fn gemini_usage(usage: &GeminiTokenUsage) -> TokenUsage {
2353    TokenUsage {
2354        input_tokens: usage.input_tokens.saturating_sub(usage.cached_tokens),
2355        cached_input_tokens: usage.cached_tokens,
2356        thinking_tokens: usage.thought_tokens,
2357        output_tokens: usage.output_tokens,
2358    }
2359}
2360
2361fn gemini_cost(cost: &kcode_gemini_api::CostBreakdown) -> CostEstimate {
2362    let accuracy = match cost.accuracy {
2363        kcode_gemini_api::CostAccuracy::Exact => CostAccuracy::Exact,
2364        kcode_gemini_api::CostAccuracy::Estimated => CostAccuracy::Estimated,
2365        kcode_gemini_api::CostAccuracy::Conservative => CostAccuracy::Conservative,
2366    };
2367    CostEstimate {
2368        usd_nanos: cost.total.usd_nanos(),
2369        accuracy,
2370        pricing_version: cost.pricing_version.clone(),
2371    }
2372}
2373
2374fn openai_image_analysis_cost(usage: &ImageAnalysisUsage) -> CostEstimate {
2375    let cached = usage.cached_input_tokens.unwrap_or(0);
2376    let cache_write = usage.cache_write_input_tokens.unwrap_or(0);
2377    let non_cached = usage
2378        .input_tokens
2379        .saturating_sub(cached)
2380        .saturating_sub(cache_write);
2381    let long = usage.input_tokens > 272_000;
2382    let (input_rate, cached_rate, cache_write_rate, output_rate) = if long {
2383        (10_000, 1_000, 12_500, 45_000)
2384    } else {
2385        (5_000, 500, 6_250, 30_000)
2386    };
2387    CostEstimate {
2388        usd_nanos: non_cached
2389            .saturating_mul(input_rate)
2390            .saturating_add(cached.saturating_mul(cached_rate))
2391            .saturating_add(cache_write.saturating_mul(cache_write_rate))
2392            .saturating_add(usage.output_tokens.saturating_mul(output_rate)),
2393        accuracy: CostAccuracy::Exact,
2394        pricing_version: PRICING_VERSION.into(),
2395    }
2396}
2397
2398fn openai_generation_cost(usage: &OpenAiGenerationUsage) -> CostEstimate {
2399    let detailed_input = usage
2400        .input_details
2401        .text_tokens
2402        .saturating_add(usage.input_details.image_tokens);
2403    let output_details = usage.output_details.as_ref();
2404    let detailed_output = output_details
2405        .map(|details| details.text_tokens.saturating_add(details.image_tokens))
2406        .unwrap_or(0);
2407    let input = usage
2408        .input_details
2409        .text_tokens
2410        .saturating_mul(5_000)
2411        .saturating_add(usage.input_details.image_tokens.saturating_mul(8_000))
2412        .saturating_add(
2413            usage
2414                .input_tokens
2415                .saturating_sub(detailed_input)
2416                .saturating_mul(8_000),
2417        );
2418    let output = detailed_output
2419        .saturating_add(usage.output_tokens.saturating_sub(detailed_output))
2420        .saturating_mul(30_000);
2421    CostEstimate {
2422        usd_nanos: input.saturating_add(output),
2423        accuracy: if detailed_input == usage.input_tokens
2424            && output_details.is_some()
2425            && detailed_output == usage.output_tokens
2426        {
2427            CostAccuracy::Exact
2428        } else {
2429            CostAccuracy::Estimated
2430        },
2431        pricing_version: PRICING_VERSION.into(),
2432    }
2433}
2434
2435fn codex_search_cost(model: &str, usage: Option<TokenUsage>) -> Option<CostEstimate> {
2436    let metering = Metering::Tokens(usage?);
2437    estimate_cost(model, &metering).map(|cost| {
2438        // The Codex search wrapper performs at least one OpenAI web-search tool call.
2439        // Provider output does not expose the exact internal search-call count.
2440        cost.with_surcharge(10_000_000, CostAccuracy::Estimated)
2441    })
2442}
2443
2444fn codex_usage(usage: &CodexTokenUsage) -> TokenUsage {
2445    TokenUsage {
2446        input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
2447        cached_input_tokens: usage.cached_input_tokens,
2448        thinking_tokens: usage.reasoning_output_tokens,
2449        output_tokens: usage
2450            .output_tokens
2451            .saturating_sub(usage.reasoning_output_tokens),
2452    }
2453}
2454
2455fn codex_v2_usage(usage: &kcode_codex_runtime_v2::TokenUsage) -> TokenUsage {
2456    TokenUsage {
2457        input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
2458        cached_input_tokens: usage.cached_input_tokens,
2459        thinking_tokens: usage.reasoning_output_tokens,
2460        output_tokens: usage
2461            .output_tokens
2462            .saturating_sub(usage.reasoning_output_tokens),
2463    }
2464}
2465
2466fn openai_image_usage(usage: &ImageAnalysisUsage) -> TokenUsage {
2467    let cached = usage.cached_input_tokens.unwrap_or(0);
2468    let thinking = usage.reasoning_output_tokens.unwrap_or(0);
2469    TokenUsage {
2470        input_tokens: usage.input_tokens.saturating_sub(cached),
2471        cached_input_tokens: cached,
2472        thinking_tokens: thinking,
2473        output_tokens: usage.output_tokens.saturating_sub(thinking),
2474    }
2475}
2476
2477fn openai_generation_usage(usage: &OpenAiGenerationUsage) -> TokenUsage {
2478    TokenUsage {
2479        input_tokens: usage.input_tokens,
2480        cached_input_tokens: 0,
2481        thinking_tokens: 0,
2482        output_tokens: usage.output_tokens,
2483    }
2484}
2485
2486fn transcription_usage(usage: TranscriptionUsage) -> Metering {
2487    match usage {
2488        TranscriptionUsage::DurationSeconds(seconds) => Metering::DurationSeconds { seconds },
2489        TranscriptionUsage::Tokens(tokens) => Metering::Tokens(TokenUsage {
2490            input_tokens: tokens.input_tokens,
2491            cached_input_tokens: 0,
2492            thinking_tokens: 0,
2493            output_tokens: tokens.output_tokens,
2494        }),
2495    }
2496}
2497
2498fn codex_error(error: kcode_codex_runtime::Error) -> Error {
2499    match error.kind() {
2500        CodexErrorKind::InvalidInput => Error::invalid(error.message()),
2501        CodexErrorKind::Authentication => {
2502            Error::unavailable("provider_not_configured", error.message())
2503        }
2504        CodexErrorKind::Unavailable => Error::unavailable("provider_unavailable", error.message()),
2505        CodexErrorKind::RateLimited | CodexErrorKind::Capacity => {
2506            Error::unavailable("provider_rate_limited", error.message())
2507        }
2508        CodexErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
2509        CodexErrorKind::InputTooLarge => Error::invalid(error.message()),
2510        CodexErrorKind::EmptyOutput | CodexErrorKind::Protocol => {
2511            Error::provider("provider_error", error.message())
2512        }
2513    }
2514}
2515
2516fn codex_v2_error(error: kcode_codex_runtime_v2::Error) -> Error {
2517    use kcode_codex_runtime_v2::ErrorKind;
2518    match error.kind() {
2519        ErrorKind::InvalidInput => Error::invalid(error.message()),
2520        ErrorKind::Unavailable | ErrorKind::Authentication => {
2521            Error::unavailable("provider_unavailable", error.message())
2522        }
2523        ErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
2524        ErrorKind::Protocol => Error::provider("provider_error", error.message()),
2525        ErrorKind::Cancelled => Error::cancelled(),
2526    }
2527}
2528
2529fn gemini_error(error: GeminiError) -> Error {
2530    match &error {
2531        GeminiError::InvalidApiKey => {
2532            Error::unavailable("provider_not_configured", error.to_string())
2533        }
2534        GeminiError::InvalidInput(_) => Error::invalid(error.to_string()),
2535        GeminiError::SpendingLimitReached { .. } => {
2536            Error::unavailable("provider_rate_limited", error.to_string())
2537        }
2538        GeminiError::Accounting(_)
2539        | GeminiError::Transport(_)
2540        | GeminiError::Provider { .. }
2541        | GeminiError::Protocol(_) => Error::provider("provider_error", error.to_string()),
2542    }
2543}
2544
2545fn openai_error(error: OpenAiError) -> Error {
2546    match &error {
2547        OpenAiError::InvalidApiKey => {
2548            Error::unavailable("provider_not_configured", error.to_string())
2549        }
2550        OpenAiError::InvalidInput(_) => Error::invalid(error.to_string()),
2551        OpenAiError::Transport(_) | OpenAiError::Provider { .. } | OpenAiError::Protocol(_) => {
2552            Error::provider("provider_error", error.to_string())
2553        }
2554    }
2555}
2556
2557fn web_fetch_error(error: kcode_web_fetch::Error) -> Error {
2558    match error.kind() {
2559        WebFetchErrorKind::InvalidInput | WebFetchErrorKind::UnsafeDestination => {
2560            Error::invalid(error.message())
2561        }
2562        WebFetchErrorKind::Timeout => Error::provider("web_fetch_timeout", error.message()),
2563        WebFetchErrorKind::UnsupportedContent => {
2564            Error::invalid(format!("unsupported web content: {}", error.message()))
2565        }
2566        WebFetchErrorKind::Transport
2567        | WebFetchErrorKind::HttpStatus
2568        | WebFetchErrorKind::EmptyContent => Error::provider("web_fetch_failed", error.message()),
2569    }
2570}
2571
2572fn document_error(error: kcode_doc_extraction::Error) -> Error {
2573    match error.kind() {
2574        DocumentErrorKind::InvalidInput | DocumentErrorKind::UnsupportedFormat => {
2575            Error::invalid(error.message())
2576        }
2577        DocumentErrorKind::ExtractionFailed | DocumentErrorKind::EmptyText => {
2578            Error::provider("document_extraction_failed", error.message())
2579        }
2580    }
2581}
2582
2583#[cfg(test)]
2584mod tests {
2585    use super::*;
2586
2587    #[test]
2588    fn audio_kind_overrides_mislabeled_ogg_video_mime() {
2589        let media = Media::audio(vec![1], "voice.ogg", "video/ogg").unwrap();
2590        assert_eq!(media.kind, MediaKind::Audio);
2591        assert_eq!(media.content_type, "audio/ogg");
2592    }
2593
2594    #[test]
2595    fn exact_search_models_replace_modes() {
2596        assert!(codex_search_profile("gpt-5.6-sol").is_ok());
2597        assert!(codex_search_profile("gpt-5.6-terra").is_ok());
2598        assert!(codex_search_profile("fast").is_err());
2599    }
2600
2601    #[test]
2602    fn long_prompt_is_accepted_and_blank_prompt_is_rejected() {
2603        assert!(validate_prompt(&"x".repeat(1_000_001)).is_ok());
2604        assert!(validate_prompt("  \n\t").is_err());
2605    }
2606
2607    #[test]
2608    fn omitted_temperature_preserves_provider_default() {
2609        let options = audio_generation_options(1_024, None).unwrap();
2610        assert_eq!(options.temperature, None);
2611    }
2612
2613    #[test]
2614    fn exact_zero_temperature_is_preserved() {
2615        let options = audio_generation_options(1_024, Some(0.0)).unwrap();
2616        assert_eq!(options.temperature, Some(0.0));
2617    }
2618
2619    #[test]
2620    fn invalid_and_nonfinite_temperatures_are_rejected_without_receipts() {
2621        for temperature in [Some(-0.1), Some(2.1), Some(f32::NAN), Some(f32::INFINITY)] {
2622            let error = validate_temperature(temperature).err().unwrap();
2623            assert!(error.receipt().is_none());
2624        }
2625    }
2626
2627    #[test]
2628    fn openai_transcription_temperature_is_rejected_without_a_receipt() {
2629        let error =
2630            validate_transcription_temperature(kcode_openai_api::GPT_4O_TRANSCRIBE, Some(0.0))
2631                .err()
2632                .unwrap();
2633        assert!(error.receipt().is_none());
2634    }
2635}