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 model_context = agent_model_context(&request, "openai");
654                let result = tokio::select! {
655                    _ = operation.cancelled() => Err(Error::cancelled()),
656                    result = tokio::time::timeout(request.timeout, openai.agent_turn(provider_request)) => {
657                        result
658                            .map_err(|_| Error::provider("provider_timeout", "OpenAI agent turn timed out."))
659                            .and_then(|result| result.map_err(openai_error))
660                    }
661                };
662                let result = self.account_result(
663                    operation_id,
664                    parent_operation_id,
665                    "agent_turn",
666                    &requested_model,
667                    result,
668                )?;
669                let actual_model = result.model.clone();
670                let provider_request_id = Some(result.response_id.clone());
671                (
672                    AgentTurnBackend::Buffered(buffered_openai_turn(
673                        provider_input,
674                        model_context,
675                        result,
676                    )),
677                    actual_model,
678                    provider_request_id,
679                )
680            }
681            AgentProvider::Gemini => {
682                let gemini = self.service.gemini.as_ref().ok_or_else(|| {
683                    Error::unavailable("provider_not_configured", "Gemini is not configured.")
684                })?;
685                let provider_request = gemini_agent_request(&request);
686                let provider_input = provider_request_json(&request);
687                let model_context = agent_model_context(&request, "gemini");
688                let result = tokio::select! {
689                    _ = operation.cancelled() => Err(Error::cancelled()),
690                    result = tokio::time::timeout(request.timeout, gemini.agent_turn(provider_request)) => {
691                        result
692                            .map_err(|_| Error::provider("provider_timeout", "Gemini agent turn timed out."))
693                            .and_then(|result| result.map_err(gemini_error))
694                    }
695                };
696                let result = self.account_result(
697                    operation_id,
698                    parent_operation_id,
699                    "agent_turn",
700                    &requested_model,
701                    result,
702                )?;
703                let actual_model = result.model.clone();
704                let provider_request_id = Some(result.interaction_id.clone());
705                (
706                    AgentTurnBackend::Buffered(buffered_gemini_turn(
707                        provider_input,
708                        model_context,
709                        result,
710                    )),
711                    actual_model,
712                    provider_request_id,
713                )
714            }
715        };
716        Ok(AgentTurn {
717            inner,
718            operation,
719            user_id: self.user_id.clone(),
720            requested_model,
721            actual_model,
722            provider_request_id,
723            operation_id,
724            parent_operation_id,
725            receipts: self.service.receipts.clone(),
726            receipt: None,
727        })
728    }
729
730    pub async fn search(&self, request: SearchRequest) -> Result<Accounted<SearchResponse>> {
731        let question = request.question.trim();
732        if question.is_empty() || question.chars().count() > 4_000 {
733            return Err(Error::invalid(
734                "question must contain between 1 and 4000 characters",
735            ));
736        }
737        validate_model(&request.model)?;
738        let mut operation = self
739            .service
740            .active_operations
741            .register_request(request.operation_id, request.parent_operation_id)
742            .map_err(registry_error)?;
743        let started = Instant::now();
744        let (response, receipt) = if let Some(model) = gemini_model(&request.model) {
745            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
746                Error::unavailable(
747                    "provider_not_configured",
748                    "Gemini search is not configured.",
749                )
750            })?;
751            let result = tokio::select! {
752                _ = operation.cancelled() => Err(Error::cancelled()),
753                result = tokio::time::timeout(
754                    FAST_SEARCH_TIMEOUT,
755                    gemini.grounded_search_with_model(
756                        model,
757                        GroundedSearchRequest::new(question),
758                    ),
759                ) => result
760                    .map_err(|_| Error::provider("provider_timeout", "Gemini search timed out."))
761                    .and_then(|result| result.map_err(gemini_error)),
762            };
763            let result = self.account_result(
764                request.operation_id,
765                request.parent_operation_id,
766                "web_search",
767                &request.model,
768                result,
769            )?;
770            let interaction = result.interaction;
771            let usage = gemini_usage(&interaction.usage);
772            let cost = gemini_cost(&interaction.cost);
773            let receipt = self.record_tokens_with_cost(
774                request.operation_id,
775                request.parent_operation_id,
776                "web_search",
777                &request.model,
778                &interaction.model,
779                Some(usage),
780                Some(interaction.id.clone()),
781                None,
782                Some(cost.clone()),
783            )?;
784            let answer = interaction
785                .text
786                .filter(|value| !value.trim().is_empty())
787                .ok_or_else(|| {
788                    Error::provider("provider_error", "Gemini search returned no answer text.")
789                })
790                .map_err(|error| error.with_receipt(receipt.clone()))?;
791            (
792                SearchResponse {
793                    answer,
794                    sources: result
795                        .sources
796                        .into_iter()
797                        .map(|source| WebSource {
798                            title: source.title,
799                            url: source.url,
800                        })
801                        .collect(),
802                    model: interaction.model,
803                    usage: Some(usage),
804                    cost: Some(cost),
805                },
806                receipt,
807            )
808        } else {
809            let (reasoning, context, depth, timeout) = codex_search_profile(&request.model)?;
810            let result = tokio::select! {
811                _ = operation.cancelled() => Err(Error::cancelled()),
812                result = self.service.codex.web_search(CodexSearchRequest {
813                    question: question.to_owned(),
814                    model: request.model.clone(),
815                    reasoning_effort: reasoning,
816                    context,
817                    depth,
818                    timeout,
819                }) => result.map_err(codex_error),
820            };
821            let result = self.account_result(
822                request.operation_id,
823                request.parent_operation_id,
824                "web_search",
825                &request.model,
826                result,
827            )?;
828            let usage = result.usage.as_ref().map(codex_usage);
829            let cost = codex_search_cost(&request.model, usage);
830            let receipt = self.record_tokens_with_cost(
831                request.operation_id,
832                request.parent_operation_id,
833                "web_search",
834                &request.model,
835                &request.model,
836                usage,
837                None,
838                None,
839                cost.clone(),
840            )?;
841            (
842                SearchResponse {
843                    answer: result.answer,
844                    sources: result
845                        .sources
846                        .into_iter()
847                        .map(|source| WebSource {
848                            title: source.title,
849                            url: source.url,
850                        })
851                        .collect(),
852                    model: request.model.clone(),
853                    usage,
854                    cost,
855                },
856                receipt,
857            )
858        };
859        tracing::info!(
860            user_id = %self.user_id,
861            operation_id = %request.operation_id,
862            model = %response.model,
863            duration_ms = started.elapsed().as_millis(),
864            "Intelligence search completed"
865        );
866        Ok(Accounted {
867            value: response,
868            receipt,
869        })
870    }
871
872    /// Performs one Gemini audio-analysis call and records its provider usage.
873    pub async fn analyze_audio(
874        &self,
875        request: AudioAnalysisRequest,
876    ) -> Result<Accounted<AudioAnalysisResponse>> {
877        validate_operation(&request.operation)?;
878        validate_prompt(&request.prompt)?;
879        validate_model(&request.model)?;
880        if request.media.kind != MediaKind::Audio {
881            return Err(Error::invalid("audio analysis requires audio media"));
882        }
883        if request.max_output_tokens == 0 || request.max_output_tokens > 65_536 {
884            return Err(Error::invalid(
885                "max_output_tokens must be between 1 and 65536",
886            ));
887        }
888        let options = audio_generation_options(request.max_output_tokens, request.temperature)?;
889        let model = gemini_model(&request.model).ok_or_else(|| {
890            Error::invalid("audio analysis requires a supported exact Gemini model")
891        })?;
892        let gemini = self.service.gemini.as_ref().ok_or_else(|| {
893            Error::unavailable(
894                "provider_not_configured",
895                "Gemini audio analysis is not configured.",
896            )
897        })?;
898        let media = gemini_media(&request.media)?;
899        let mut provider_request = MultimodalRequest::new(request.prompt, vec![media]);
900        provider_request.options = options;
901        provider_request.structured_output = request
902            .schema
903            .map(StructuredOutput::new)
904            .transpose()
905            .map_err(gemini_error)?;
906        let mut operation = self
907            .service
908            .active_operations
909            .register_request(request.operation_id, request.parent_operation_id)
910            .map_err(registry_error)?;
911        let result = tokio::select! {
912            _ = operation.cancelled() => Err(Error::cancelled()),
913            result = tokio::time::timeout(
914                MEDIA_ANNOTATION_TIMEOUT,
915                gemini.infer_multimodal(model, provider_request),
916            ) => result
917                .map_err(|_| Error::provider("provider_timeout", "Gemini audio analysis timed out."))
918                .and_then(|result| result.map_err(gemini_error)),
919        };
920        let result = self.account_result(
921            request.operation_id,
922            request.parent_operation_id,
923            &request.operation,
924            &request.model,
925            result,
926        )?;
927        let usage = gemini_usage(&result.usage);
928        let cost = gemini_cost(&result.cost);
929        let receipt = self.record_tokens_with_cost(
930            request.operation_id,
931            request.parent_operation_id,
932            &request.operation,
933            &request.model,
934            &result.model,
935            Some(usage),
936            Some(result.id.clone()),
937            None,
938            Some(cost.clone()),
939        )?;
940        if result.status != GeminiCompletionStatus::Completed {
941            return Err(Error::provider(
942                "provider_incomplete",
943                "Gemini audio analysis did not complete.",
944            )
945            .with_receipt(receipt));
946        }
947        let text = result
948            .text
949            .filter(|text| !text.trim().is_empty())
950            .ok_or_else(|| {
951                Error::provider(
952                    "provider_empty_output",
953                    "Gemini audio analysis returned no text.",
954                )
955            })
956            .map_err(|error| error.with_receipt(receipt.clone()))?;
957        Ok(Accounted {
958            value: AudioAnalysisResponse {
959                model: result.model,
960                text,
961                usage,
962                cost,
963            },
964            receipt,
965        })
966    }
967
968    /// Performs one structured Gemini audio call and records its provider usage.
969    pub async fn transcribe_structured_audio(
970        &self,
971        request: StructuredAudioRequest,
972    ) -> Result<Accounted<StructuredAudioResponse>> {
973        self.analyze_audio(AudioAnalysisRequest {
974            operation: request.operation,
975            prompt: request.prompt,
976            model: request.model,
977            media: request.media,
978            schema: Some(request.schema),
979            max_output_tokens: request.max_output_tokens,
980            temperature: request.temperature,
981            operation_id: request.operation_id,
982            parent_operation_id: request.parent_operation_id,
983        })
984        .await
985        .map(|response| {
986            response.map(|response| StructuredAudioResponse {
987                model: response.model,
988                text: response.text,
989                usage: response.usage,
990                cost: response.cost,
991            })
992        })
993    }
994
995    /// Performs one tool-free Codex generation and records its provider usage.
996    pub async fn generate_text(
997        &self,
998        request: TextGenerationRequest,
999    ) -> Result<Accounted<TextGenerationResponse>> {
1000        validate_operation(&request.operation)?;
1001        validate_model(&request.model)?;
1002        if request.prompt.trim().is_empty() || request.prompt.chars().count() > 1_000_000 {
1003            return Err(Error::invalid(
1004                "generation prompt must contain between 1 and 1000000 characters",
1005            ));
1006        }
1007        if request.timeout.is_zero() || request.timeout > Duration::from_secs(7 * 60 * 60 + 30 * 60)
1008        {
1009            return Err(Error::invalid(
1010                "generation timeout must be between 1 second and 7 hours 30 minutes",
1011            ));
1012        }
1013        let mut operation = self
1014            .service
1015            .active_operations
1016            .register_request(request.operation_id, request.parent_operation_id)
1017            .map_err(registry_error)?;
1018        let mut provider_request =
1019            CodexGenerationRequest::new(request.prompt, request.model.clone());
1020        provider_request.reasoning_effort = request.reasoning_effort;
1021        provider_request.ephemeral = true;
1022        provider_request.timeout = request.timeout;
1023        let result = tokio::select! {
1024            _ = operation.cancelled() => Err(Error::cancelled()),
1025            result = self.service.codex.generate(provider_request) => result.map_err(codex_error),
1026        };
1027        let result = self.account_result(
1028            request.operation_id,
1029            request.parent_operation_id,
1030            &request.operation,
1031            &request.model,
1032            result,
1033        )?;
1034        let usage = result.usage.as_ref().map(codex_usage);
1035        let cost = usage.and_then(|usage| estimate_token_cost(&request.model, usage));
1036        let receipt = self.record_tokens(
1037            request.operation_id,
1038            request.parent_operation_id,
1039            &request.operation,
1040            &request.model,
1041            &request.model,
1042            usage,
1043            None,
1044            Some(result.thread_id.clone()),
1045        )?;
1046        Ok(Accounted {
1047            value: TextGenerationResponse {
1048                model: request.model,
1049                text: result.answer,
1050                thread_id: result.thread_id,
1051                usage,
1052                cost,
1053            },
1054            receipt,
1055        })
1056    }
1057
1058    pub async fn fetch(&self, request: FetchRequest) -> Result<FetchResponse> {
1059        let mut operation = self
1060            .service
1061            .active_operations
1062            .register_request(request.operation_id, request.parent_operation_id)
1063            .map_err(registry_error)?;
1064        let fetched = tokio::select! {
1065            _ = operation.cancelled() => Err(Error::cancelled()),
1066            result = self.service.web_fetcher.fetch(&request.url) => result.map_err(web_fetch_error),
1067        }?;
1068        Ok(FetchResponse {
1069            url: fetched.url,
1070            title: fetched.title,
1071            content_type: fetched.content_type,
1072            content: fetched.content,
1073            truncated: fetched.truncated,
1074            retrieved_at: DateTime::<Utc>::from(fetched.retrieved_at),
1075        })
1076    }
1077
1078    pub async fn transcribe(
1079        &self,
1080        request: TranscriptionRequest,
1081    ) -> Result<Accounted<TranscriptionResponse>> {
1082        validate_prompt(&request.prompt)?;
1083        if request.media.kind != MediaKind::Audio {
1084            return Err(Error::invalid("transcription requires audio media"));
1085        }
1086        validate_model(&request.model)?;
1087        let temperature = validate_transcription_temperature(&request.model, request.temperature)?;
1088        let mut operation = self
1089            .service
1090            .active_operations
1091            .register_request(request.operation_id, request.parent_operation_id)
1092            .map_err(registry_error)?;
1093        let (response, receipt) = if request.model == kcode_openai_api::GPT_4O_TRANSCRIBE {
1094            let openai = self.service.openai.as_ref().ok_or_else(|| {
1095                Error::unavailable(
1096                    "transcription_unavailable",
1097                    "OpenAI audio transcription is not configured.",
1098                )
1099            })?;
1100            let input = AudioInput::new(
1101                safe_audio_filename(&request.media.file_name, &request.media.content_type),
1102                request.media.content_type.clone(),
1103                request.media.bytes,
1104            )
1105            .map_err(openai_error)?;
1106            let mut provider_request = OpenAiTranscriptionRequest::new(input);
1107            provider_request.prompt = Some(request.prompt);
1108            let result = tokio::select! {
1109                _ = operation.cancelled() => Err(Error::cancelled()),
1110                result = tokio::time::timeout(
1111                    MEDIA_ANNOTATION_TIMEOUT,
1112                    openai.transcribe(provider_request),
1113                ) => result
1114                    .map_err(|_| Error::provider("provider_timeout", "OpenAI transcription timed out."))
1115                    .and_then(|result| result.map_err(openai_error)),
1116            };
1117            let result = self.account_result(
1118                request.operation_id,
1119                request.parent_operation_id,
1120                "transcribe_audio",
1121                &request.model,
1122                result,
1123            )?;
1124            let metering = result
1125                .usage
1126                .map(transcription_usage)
1127                .unwrap_or(Metering::Unavailable);
1128            let cost = estimate_cost(&request.model, &metering);
1129            let receipt = self.record_metering(
1130                request.operation_id,
1131                request.parent_operation_id,
1132                "transcribe_audio",
1133                &request.model,
1134                &request.model,
1135                metering.clone(),
1136                None,
1137                None,
1138            )?;
1139            (
1140                TranscriptionResponse {
1141                    model: request.model.clone(),
1142                    text: result.text,
1143                    metering,
1144                    cost,
1145                },
1146                receipt,
1147            )
1148        } else {
1149            let model = gemini_model(&request.model).ok_or_else(|| {
1150                Error::invalid("transcription model must be gpt-4o-transcribe or a supported exact Gemini model")
1151            })?;
1152            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1153                Error::unavailable(
1154                    "provider_not_configured",
1155                    "Gemini transcription is not configured.",
1156                )
1157            })?;
1158            let media = gemini_media(&request.media)?;
1159            let mut provider_request = MultimodalRequest::new(request.prompt, vec![media]);
1160            provider_request.options.temperature = temperature;
1161            let result = tokio::select! {
1162                _ = operation.cancelled() => Err(Error::cancelled()),
1163                result = tokio::time::timeout(
1164                    MEDIA_ANNOTATION_TIMEOUT,
1165                    gemini.infer_multimodal(model, provider_request),
1166                ) => result
1167                    .map_err(|_| Error::provider("provider_timeout", "Gemini transcription timed out."))
1168                    .and_then(|result| result.map_err(gemini_error)),
1169            };
1170            let result = self.account_result(
1171                request.operation_id,
1172                request.parent_operation_id,
1173                "transcribe_audio",
1174                &request.model,
1175                result,
1176            )?;
1177            let metering = Metering::Tokens(gemini_usage(&result.usage));
1178            let cost = gemini_cost(&result.cost);
1179            let receipt = self.record_metering_with_cost(
1180                request.operation_id,
1181                request.parent_operation_id,
1182                "transcribe_audio",
1183                &request.model,
1184                &result.model,
1185                metering.clone(),
1186                Some(result.id.clone()),
1187                None,
1188                Some(cost.clone()),
1189            )?;
1190            let text = result
1191                .text
1192                .filter(|text| !text.trim().is_empty())
1193                .ok_or_else(|| {
1194                    Error::provider(
1195                        "empty_transcription",
1196                        "Gemini returned no transcription text.",
1197                    )
1198                })
1199                .map_err(|error| error.with_receipt(receipt.clone()))?;
1200            (
1201                TranscriptionResponse {
1202                    model: result.model,
1203                    text,
1204                    metering,
1205                    cost: Some(cost),
1206                },
1207                receipt,
1208            )
1209        };
1210        Ok(Accounted {
1211            value: response,
1212            receipt,
1213        })
1214    }
1215
1216    pub async fn annotate(
1217        &self,
1218        request: AnnotationRequest,
1219    ) -> Result<Accounted<AnnotationResponse>> {
1220        validate_prompt(&request.prompt)?;
1221        validate_model(&request.model)?;
1222        let mut operation = self
1223            .service
1224            .active_operations
1225            .register_request(request.operation_id, request.parent_operation_id)
1226            .map_err(registry_error)?;
1227        let file_name = request.media.file_name.clone();
1228        let content_type = request.media.content_type.clone();
1229        let (response, receipt) = if let Some(model) = gemini_model(&request.model) {
1230            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1231                Error::unavailable(
1232                    "provider_not_configured",
1233                    "Gemini media annotation is not configured.",
1234                )
1235            })?;
1236            let media = gemini_media(&request.media)?;
1237            let result = tokio::select! {
1238                _ = operation.cancelled() => Err(Error::cancelled()),
1239                result = tokio::time::timeout(
1240                    MEDIA_ANNOTATION_TIMEOUT,
1241                    gemini.infer_multimodal(
1242                        model,
1243                        MultimodalRequest::new(request.prompt, vec![media]),
1244                    ),
1245                ) => result
1246                    .map_err(|_| Error::provider("provider_timeout", "Gemini annotation timed out."))
1247                    .and_then(|result| result.map_err(gemini_error)),
1248            };
1249            let result = self.account_result(
1250                request.operation_id,
1251                request.parent_operation_id,
1252                "annotate_media",
1253                &request.model,
1254                result,
1255            )?;
1256            let usage = gemini_usage(&result.usage);
1257            let cost = gemini_cost(&result.cost);
1258            let receipt = self.record_tokens_with_cost(
1259                request.operation_id,
1260                request.parent_operation_id,
1261                "annotate_media",
1262                &request.model,
1263                &result.model,
1264                Some(usage),
1265                Some(result.id.clone()),
1266                None,
1267                Some(cost.clone()),
1268            )?;
1269            let text = result
1270                .text
1271                .filter(|text| !text.trim().is_empty())
1272                .ok_or_else(|| {
1273                    Error::provider("empty_annotation", "Gemini returned no annotation text.")
1274                })
1275                .map_err(|error| error.with_receipt(receipt.clone()))?;
1276            (
1277                AnnotationResponse {
1278                    complete: result.status == GeminiCompletionStatus::Completed,
1279                    model: result.model,
1280                    file_name,
1281                    content_type,
1282                    text,
1283                    incomplete_reason: None,
1284                    usage: Some(usage),
1285                    cost: Some(cost),
1286                },
1287                receipt,
1288            )
1289        } else if request.model == "gpt-5.6" {
1290            if request.media.kind != MediaKind::Image {
1291                return Err(Error::invalid(
1292                    "OpenAI media annotation accepts images only",
1293                ));
1294            }
1295            let openai = self.service.openai.as_ref().ok_or_else(|| {
1296                Error::unavailable(
1297                    "provider_not_configured",
1298                    "OpenAI media annotation is not configured.",
1299                )
1300            })?;
1301            let image = OpenAiImageInput::new(
1302                openai_image_media_type(&request.media.content_type)?,
1303                request.media.bytes,
1304            )
1305            .map_err(openai_error)?;
1306            let result = tokio::select! {
1307                _ = operation.cancelled() => Err(Error::cancelled()),
1308                result = tokio::time::timeout(
1309                    MEDIA_ANNOTATION_TIMEOUT,
1310                    openai.analyze_image(ImageAnalysisRequest::new(image, request.prompt)),
1311                ) => result
1312                    .map_err(|_| Error::provider("provider_timeout", "OpenAI annotation timed out."))
1313                    .and_then(|result| result.map_err(openai_error)),
1314            };
1315            let result = self.account_result(
1316                request.operation_id,
1317                request.parent_operation_id,
1318                "annotate_media",
1319                &request.model,
1320                result,
1321            )?;
1322            let usage = result.usage.as_ref().map(openai_image_usage);
1323            let cost = result.usage.as_ref().map(openai_image_analysis_cost);
1324            let receipt = self.record_tokens_with_cost(
1325                request.operation_id,
1326                request.parent_operation_id,
1327                "annotate_media",
1328                &request.model,
1329                &result.model,
1330                usage,
1331                None,
1332                None,
1333                cost.clone(),
1334            )?;
1335            let (complete, incomplete_reason) = match result.status {
1336                OpenAiImageStatus::Completed => (true, None),
1337                OpenAiImageStatus::Incomplete { reason } => (false, reason),
1338            };
1339            (
1340                AnnotationResponse {
1341                    complete,
1342                    model: result.model,
1343                    file_name,
1344                    content_type,
1345                    text: result.text,
1346                    incomplete_reason,
1347                    usage,
1348                    cost,
1349                },
1350                receipt,
1351            )
1352        } else if matches!(
1353            request.model.as_str(),
1354            "gpt-5.6-sol" | "gpt-5.6-terra" | "gpt-5.6-luna"
1355        ) {
1356            if request.media.kind != MediaKind::Image {
1357                return Err(Error::invalid("Codex media annotation accepts images only"));
1358            }
1359            let image = kcode_codex_runtime_v2::ImageInput::new(
1360                codex_image_media_type(&request.media.content_type)?,
1361                request.media.bytes,
1362            )
1363            .map_err(codex_v2_error)?;
1364            let turn_result = self
1365                .service
1366                .agent
1367                .start_image_turn(kcode_codex_runtime_v2::ImageTurnRequest::new(
1368                    request.prompt,
1369                    request.model.clone(),
1370                    vec![image],
1371                ))
1372                .await
1373                .map_err(codex_v2_error);
1374            let mut turn = self.account_result(
1375                request.operation_id,
1376                request.parent_operation_id,
1377                "annotate_media",
1378                &request.model,
1379                turn_result,
1380            )?;
1381            let completed = loop {
1382                let event = tokio::select! {
1383                    _ = operation.cancelled() => {
1384                        turn.cancel();
1385                        let receipt = self.record_metering(
1386                            request.operation_id,
1387                            request.parent_operation_id,
1388                            "annotate_media",
1389                            &request.model,
1390                            &request.model,
1391                            Metering::Unavailable,
1392                            None,
1393                            None,
1394                        )?;
1395                        return Err(Error::cancelled().with_receipt(receipt));
1396                    }
1397                    event = turn.next_event() => event,
1398                };
1399                match event {
1400                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::ProviderInput(_))) => {}
1401                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::ModelContextSubmitted(_))) => {}
1402                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::UsageUpdated(_))) => {}
1403                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::ToolCall(_))) => {
1404                        turn.cancel();
1405                        let receipt = self.record_metering(
1406                            request.operation_id,
1407                            request.parent_operation_id,
1408                            "annotate_media",
1409                            &request.model,
1410                            &request.model,
1411                            Metering::Unavailable,
1412                            None,
1413                            None,
1414                        )?;
1415                        return Err(Error::provider(
1416                            "provider_error",
1417                            "A tool-free Codex image turn requested a tool.",
1418                        )
1419                        .with_receipt(receipt));
1420                    }
1421                    Some(Ok(kcode_codex_runtime_v2::AgentEvent::Completed(completed))) => {
1422                        break completed;
1423                    }
1424                    Some(Err(error)) => {
1425                        let receipt = self.record_metering(
1426                            request.operation_id,
1427                            request.parent_operation_id,
1428                            "annotate_media",
1429                            &request.model,
1430                            &request.model,
1431                            Metering::Unavailable,
1432                            None,
1433                            None,
1434                        )?;
1435                        return Err(codex_v2_error(error).with_receipt(receipt));
1436                    }
1437                    None => {
1438                        let receipt = self.record_metering(
1439                            request.operation_id,
1440                            request.parent_operation_id,
1441                            "annotate_media",
1442                            &request.model,
1443                            &request.model,
1444                            Metering::Unavailable,
1445                            None,
1446                            None,
1447                        )?;
1448                        return Err(Error::provider(
1449                            "empty_annotation",
1450                            "Codex ended without annotation text.",
1451                        )
1452                        .with_receipt(receipt));
1453                    }
1454                }
1455            };
1456            let usage = completed.usage.as_ref().map(codex_v2_usage);
1457            let cost = usage.and_then(|usage| estimate_token_cost(&request.model, usage));
1458            let receipt = self.record_tokens(
1459                request.operation_id,
1460                request.parent_operation_id,
1461                "annotate_media",
1462                &request.model,
1463                &request.model,
1464                usage,
1465                Some(completed.turn_id.clone()),
1466                Some(completed.thread_id.clone()),
1467            )?;
1468            (
1469                AnnotationResponse {
1470                    complete: true,
1471                    model: request.model.clone(),
1472                    file_name,
1473                    content_type,
1474                    text: completed.answer,
1475                    incomplete_reason: None,
1476                    usage,
1477                    cost,
1478                },
1479                receipt,
1480            )
1481        } else {
1482            return Err(Error::invalid(format!(
1483                "unsupported exact annotation model {}",
1484                request.model
1485            )));
1486        };
1487        Ok(Accounted {
1488            value: response,
1489            receipt,
1490        })
1491    }
1492
1493    /// Creates a new image or modifies supplied reference images.
1494    pub async fn generate_image(&self, request: ImageRequest) -> Result<Accounted<ImageResponse>> {
1495        validate_model(&request.model)?;
1496        if request.prompt.trim().is_empty()
1497            || request.prompt.chars().count() > MAX_IMAGE_PROMPT_CHARACTERS
1498        {
1499            return Err(Error::invalid(format!(
1500                "image prompt must contain 1 through {MAX_IMAGE_PROMPT_CHARACTERS} characters"
1501            )));
1502        }
1503        if request
1504            .references
1505            .iter()
1506            .any(|media| media.kind != MediaKind::Image)
1507        {
1508            return Err(Error::invalid("image references must all be images"));
1509        }
1510        let mut operation = self
1511            .service
1512            .active_operations
1513            .register_request(request.operation_id, request.parent_operation_id)
1514            .map_err(registry_error)?;
1515        let operation_name = if request.references.is_empty() {
1516            "generate_image"
1517        } else {
1518            "edit_image"
1519        };
1520        if request.model == kcode_openai_api::GPT_IMAGE_2 {
1521            let openai = self.service.openai.as_ref().ok_or_else(|| {
1522                Error::unavailable(
1523                    "provider_not_configured",
1524                    "OpenAI image generation is not configured.",
1525                )
1526            })?;
1527            let result = if request.references.is_empty() {
1528                let provider_request = OpenAiImageRequest::new(request.prompt);
1529                tokio::select! {
1530                    _ = operation.cancelled() => Err(Error::cancelled()),
1531                    result = tokio::time::timeout(
1532                        IMAGE_OPERATION_TIMEOUT,
1533                        openai.generate_image(provider_request),
1534                    ) => result
1535                        .map_err(|_| Error::provider("provider_timeout", "OpenAI image generation timed out."))
1536                        .and_then(|result| result.map_err(openai_error)),
1537                }
1538            } else {
1539                let mut images = request
1540                    .references
1541                    .into_iter()
1542                    .map(|media| {
1543                        OpenAiImageInput::new(
1544                            openai_image_media_type(&media.content_type)?,
1545                            media.bytes,
1546                        )
1547                        .map_err(openai_error)
1548                    })
1549                    .collect::<Result<Vec<_>>>()?;
1550                let first = images.remove(0);
1551                let mut provider_request = OpenAiImageEditRequest::new(first, request.prompt);
1552                provider_request.images.extend(images);
1553                tokio::select! {
1554                    _ = operation.cancelled() => Err(Error::cancelled()),
1555                    result = tokio::time::timeout(
1556                        IMAGE_OPERATION_TIMEOUT,
1557                        openai.edit_image(provider_request),
1558                    ) => result
1559                        .map_err(|_| Error::provider("provider_timeout", "OpenAI image editing timed out."))
1560                        .and_then(|result| result.map_err(openai_error)),
1561                }
1562            };
1563            let result = self.account_result(
1564                request.operation_id,
1565                request.parent_operation_id,
1566                operation_name,
1567                &request.model,
1568                result,
1569            )?;
1570            let usage = result.usage.as_ref().map(openai_generation_usage);
1571            let cost = result.usage.as_ref().map(openai_generation_cost);
1572            let receipt = self.record_tokens_with_cost(
1573                request.operation_id,
1574                request.parent_operation_id,
1575                operation_name,
1576                &request.model,
1577                kcode_openai_api::GPT_IMAGE_2,
1578                usage,
1579                result.request_id,
1580                None,
1581                cost.clone(),
1582            )?;
1583            Ok(Accounted {
1584                value: ImageResponse {
1585                    model: kcode_openai_api::GPT_IMAGE_2.into(),
1586                    content_type: result.image.format.mime_type().into(),
1587                    bytes: result.image.data,
1588                    usage,
1589                    cost,
1590                },
1591                receipt,
1592            })
1593        } else if request.model == kcode_gemini_api::NANO_BANANA_PRO {
1594            let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1595                Error::unavailable(
1596                    "provider_not_configured",
1597                    "Gemini image generation is not configured.",
1598                )
1599            })?;
1600            let mut provider_request = NanoBananaProRequest::new(request.prompt);
1601            provider_request.images = request
1602                .references
1603                .into_iter()
1604                .map(|media| {
1605                    GeminiMediaInput::image(&media.content_type, media.bytes).map_err(gemini_error)
1606                })
1607                .collect::<Result<Vec<_>>>()?;
1608            let result = tokio::select! {
1609                _ = operation.cancelled() => Err(Error::cancelled()),
1610                result = tokio::time::timeout(
1611                    IMAGE_OPERATION_TIMEOUT,
1612                    gemini.nano_banana_pro(provider_request),
1613                ) => result
1614                    .map_err(|_| Error::provider("provider_timeout", "Gemini image generation timed out."))
1615                    .and_then(|result| result.map_err(gemini_error)),
1616            };
1617            let result = self.account_result(
1618                request.operation_id,
1619                request.parent_operation_id,
1620                operation_name,
1621                &request.model,
1622                result,
1623            )?;
1624            let usage = gemini_usage(&result.usage);
1625            let cost = gemini_cost(&result.cost);
1626            let receipt = self.record_tokens_with_cost(
1627                request.operation_id,
1628                request.parent_operation_id,
1629                operation_name,
1630                &request.model,
1631                &result.model,
1632                Some(usage),
1633                Some(result.id.clone()),
1634                None,
1635                Some(cost.clone()),
1636            )?;
1637            let mut images = result.images;
1638            if images.len() != 1 {
1639                return Err(Error::provider(
1640                    "provider_error",
1641                    "Gemini did not return exactly one generated image.",
1642                )
1643                .with_receipt(receipt));
1644            }
1645            let image = images.remove(0);
1646            Ok(Accounted {
1647                value: ImageResponse {
1648                    model: result.model,
1649                    content_type: image.mime_type,
1650                    bytes: image.data,
1651                    usage: Some(usage),
1652                    cost: Some(cost),
1653                },
1654                receipt,
1655            })
1656        } else {
1657            Err(Error::invalid(format!(
1658                "unsupported exact image model {}",
1659                request.model
1660            )))
1661        }
1662    }
1663
1664    #[allow(clippy::too_many_arguments)]
1665    fn record_tokens(
1666        &self,
1667        operation_id: Uuid,
1668        parent_operation_id: Option<Uuid>,
1669        operation: &str,
1670        requested_model: &str,
1671        actual_model: &str,
1672        usage: Option<TokenUsage>,
1673        provider_request_id: Option<String>,
1674        provider_thread_id: Option<String>,
1675    ) -> Result<UsageReceipt> {
1676        self.record_metering_with_cost(
1677            operation_id,
1678            parent_operation_id,
1679            operation,
1680            requested_model,
1681            actual_model,
1682            usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
1683            provider_request_id,
1684            provider_thread_id,
1685            None,
1686        )
1687    }
1688
1689    #[allow(clippy::too_many_arguments)]
1690    fn record_tokens_with_cost(
1691        &self,
1692        operation_id: Uuid,
1693        parent_operation_id: Option<Uuid>,
1694        operation: &str,
1695        requested_model: &str,
1696        actual_model: &str,
1697        usage: Option<TokenUsage>,
1698        provider_request_id: Option<String>,
1699        provider_thread_id: Option<String>,
1700        cost: Option<CostEstimate>,
1701    ) -> Result<UsageReceipt> {
1702        self.record_metering_with_cost(
1703            operation_id,
1704            parent_operation_id,
1705            operation,
1706            requested_model,
1707            actual_model,
1708            usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
1709            provider_request_id,
1710            provider_thread_id,
1711            cost,
1712        )
1713    }
1714
1715    #[allow(clippy::too_many_arguments)]
1716    fn record_metering(
1717        &self,
1718        operation_id: Uuid,
1719        parent_operation_id: Option<Uuid>,
1720        operation: &str,
1721        requested_model: &str,
1722        actual_model: &str,
1723        metering: Metering,
1724        provider_request_id: Option<String>,
1725        provider_thread_id: Option<String>,
1726    ) -> Result<UsageReceipt> {
1727        self.record_metering_with_cost(
1728            operation_id,
1729            parent_operation_id,
1730            operation,
1731            requested_model,
1732            actual_model,
1733            metering,
1734            provider_request_id,
1735            provider_thread_id,
1736            None,
1737        )
1738    }
1739
1740    #[allow(clippy::too_many_arguments)]
1741    fn record_metering_with_cost(
1742        &self,
1743        operation_id: Uuid,
1744        parent_operation_id: Option<Uuid>,
1745        operation: &str,
1746        requested_model: &str,
1747        actual_model: &str,
1748        metering: Metering,
1749        provider_request_id: Option<String>,
1750        provider_thread_id: Option<String>,
1751        cost: Option<CostEstimate>,
1752    ) -> Result<UsageReceipt> {
1753        let mut receipt = UsageReceipt::new(
1754            self.user_id.clone(),
1755            operation_id,
1756            parent_operation_id,
1757            operation,
1758            requested_model,
1759            actual_model,
1760            metering,
1761        );
1762        receipt.provider_request_id = provider_request_id;
1763        receipt.provider_thread_id = provider_thread_id;
1764        if cost.is_some() {
1765            receipt.cost = cost;
1766        }
1767        self.service.receipts.record(&receipt)?;
1768        Ok(receipt)
1769    }
1770
1771    fn account_result<T>(
1772        &self,
1773        operation_id: Uuid,
1774        parent_operation_id: Option<Uuid>,
1775        operation: &str,
1776        requested_model: &str,
1777        result: Result<T>,
1778    ) -> Result<T> {
1779        match result {
1780            Ok(value) => Ok(value),
1781            Err(error) => {
1782                let receipt = self.record_metering(
1783                    operation_id,
1784                    parent_operation_id,
1785                    operation,
1786                    requested_model,
1787                    requested_model,
1788                    Metering::Unavailable,
1789                    None,
1790                    None,
1791                )?;
1792                Err(error.with_receipt(receipt))
1793            }
1794        }
1795    }
1796}
1797
1798async fn extract_document(
1799    extractor: &DocumentExtractor,
1800    document: Document,
1801) -> Result<DocumentExtraction> {
1802    let extractor = extractor.clone();
1803    let extracted = tokio::task::spawn_blocking(move || {
1804        extractor.extract(DocumentInput {
1805            file_name: document.file_name,
1806            content_type: document.content_type,
1807            data: document.bytes,
1808        })
1809    })
1810    .await
1811    .map_err(|_| {
1812        Error::internal(
1813            "document_extraction_failed",
1814            "The document extraction worker stopped unexpectedly.",
1815        )
1816    })?
1817    .map_err(document_error)?;
1818    Ok(DocumentExtraction {
1819        file_name: extracted.file_name,
1820        content_type: extracted.content_type,
1821        format: extracted.format.as_str().into(),
1822        text: extracted.text,
1823        characters: extracted.characters,
1824        truncated: extracted.truncated,
1825    })
1826}
1827
1828impl AgentTurn {
1829    /// Canonical receipt after the turn has completed or been marked unavailable.
1830    pub fn receipt(&self) -> Option<&UsageReceipt> {
1831        self.receipt.as_ref()
1832    }
1833
1834    /// Ends an interrupted turn with unavailable metering and returns its receipt.
1835    pub fn finish_unavailable(&mut self) -> Result<&UsageReceipt> {
1836        if self.receipt.is_none() {
1837            self.record_unavailable()?;
1838        }
1839        Ok(self.receipt.as_ref().expect("receipt was just recorded"))
1840    }
1841
1842    pub async fn next_event(&mut self) -> Result<Option<kcode_codex_runtime_v2::AgentEvent>> {
1843        let event = match &mut self.inner {
1844            AgentTurnBackend::Codex(inner) => tokio::select! {
1845                _ = self.operation.cancelled() => {
1846                    inner.cancel();
1847                    let receipt = self.record_unavailable()?;
1848                    return Err(Error::cancelled().with_receipt(receipt));
1849                }
1850                event = inner.next_event() => match event {
1851                    Some(Ok(event)) => Some(event),
1852                    Some(Err(error)) => {
1853                        let receipt = self.record_unavailable()?;
1854                        return Err(codex_v2_error(error).with_receipt(receipt));
1855                    }
1856                    None => {
1857                        self.record_unavailable()?;
1858                        None
1859                    },
1860                }
1861            },
1862            AgentTurnBackend::Buffered(inner) => {
1863                if self.operation.direct_cancellation_requested() {
1864                    let receipt = self.record_unavailable()?;
1865                    return Err(Error::cancelled().with_receipt(receipt));
1866                }
1867                inner.events.pop_front()
1868            }
1869        };
1870        if let Some(kcode_codex_runtime_v2::AgentEvent::Completed(completed)) = &event
1871            && self.receipt.is_none()
1872        {
1873            let mut receipt = UsageReceipt::new(
1874                self.user_id.clone(),
1875                self.operation_id,
1876                self.parent_operation_id,
1877                "agent_turn",
1878                self.requested_model.clone(),
1879                self.actual_model.clone(),
1880                completed
1881                    .usage
1882                    .as_ref()
1883                    .map(codex_v2_usage)
1884                    .map(Metering::Tokens)
1885                    .unwrap_or(Metering::Unavailable),
1886            );
1887            receipt.provider_request_id = self
1888                .provider_request_id
1889                .clone()
1890                .or_else(|| Some(completed.turn_id.clone()));
1891            receipt.provider_thread_id = Some(completed.thread_id.clone());
1892            self.receipts.record(&receipt)?;
1893            self.receipt = Some(receipt);
1894        }
1895        Ok(event)
1896    }
1897
1898    fn record_unavailable(&mut self) -> Result<UsageReceipt> {
1899        if let Some(receipt) = &self.receipt {
1900            return Ok(receipt.clone());
1901        }
1902        let receipt = UsageReceipt::new(
1903            self.user_id.clone(),
1904            self.operation_id,
1905            self.parent_operation_id,
1906            "agent_turn",
1907            self.requested_model.clone(),
1908            self.actual_model.clone(),
1909            Metering::Unavailable,
1910        );
1911        self.receipts.record(&receipt)?;
1912        self.receipt = Some(receipt.clone());
1913        Ok(receipt)
1914    }
1915
1916    pub async fn respond(
1917        &mut self,
1918        call_id: &str,
1919        result: kcode_codex_runtime_v2::ToolResult,
1920    ) -> Result<()> {
1921        match &mut self.inner {
1922            AgentTurnBackend::Codex(inner) => {
1923                inner.respond(call_id, result).await.map_err(codex_v2_error)
1924            }
1925            AgentTurnBackend::Buffered(inner) => {
1926                if inner.pending_call_id.as_deref() != Some(call_id) {
1927                    return Err(Error::invalid(
1928                        "tool result does not match the pending provider call",
1929                    ));
1930                }
1931                inner.pending_call_id = None;
1932                if let Some(mut completed) = inner.completed.take() {
1933                    completed.answer.clear();
1934                    inner
1935                        .events
1936                        .push_back(kcode_codex_runtime_v2::AgentEvent::Completed(completed));
1937                }
1938                Ok(())
1939            }
1940        }
1941    }
1942}
1943
1944impl Drop for AgentTurn {
1945    fn drop(&mut self) {
1946        let _ = self.record_unavailable();
1947    }
1948}
1949
1950fn registry_error(error: RegistryError) -> Error {
1951    match error {
1952        RegistryError::RegistryUnavailable => Error::internal(
1953            "operation_registry_unavailable",
1954            "The operation registry is unavailable.",
1955        ),
1956        RegistryError::SameAsParent => {
1957            Error::invalid("operation_id and parent_operation_id must be different")
1958        }
1959        RegistryError::OperationInProgress => Error::conflict(
1960            "operation_in_progress",
1961            "An operation with this identifier is already running.",
1962        ),
1963        RegistryError::ParentNotRunning => Error::conflict(
1964            "parent_operation_not_running",
1965            "The parent operation is no longer running.",
1966        ),
1967    }
1968}
1969
1970fn validate_model(model: &str) -> Result<()> {
1971    if model.trim().is_empty()
1972        || model.chars().count() > 128
1973        || !model
1974            .bytes()
1975            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_'))
1976    {
1977        return Err(Error::invalid(
1978            "model must be an exact safe model identifier",
1979        ));
1980    }
1981    Ok(())
1982}
1983
1984fn validate_agent_model(model: &str) -> Result<()> {
1985    if model.trim().is_empty()
1986        || model.chars().count() > 128
1987        || !model.bytes().all(|byte| {
1988            byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_' | b'/' | b':')
1989        })
1990        || model.starts_with("codex:")
1991    {
1992        return Err(Error::invalid(
1993            "model must be an exact safe provider model identifier",
1994        ));
1995    }
1996    Ok(())
1997}
1998
1999fn validate_operation(operation: &str) -> Result<()> {
2000    if operation.is_empty()
2001        || operation.len() > 64
2002        || !operation
2003            .bytes()
2004            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
2005    {
2006        return Err(Error::invalid(
2007            "operation must be a lowercase identifier of at most 64 bytes",
2008        ));
2009    }
2010    Ok(())
2011}
2012
2013fn validate_prompt(prompt: &str) -> Result<()> {
2014    if prompt.trim().is_empty() {
2015        return Err(Error::invalid("prompt must not be blank"));
2016    }
2017    Ok(())
2018}
2019
2020fn validate_temperature(temperature: Option<f32>) -> Result<Option<f32>> {
2021    match temperature {
2022        Some(value) if value.is_finite() && (0.0..=2.0).contains(&value) => Ok(Some(value)),
2023        Some(_) => Err(Error::invalid(
2024            "temperature must be finite and between 0.0 and 2.0 inclusive",
2025        )),
2026        None => Ok(None),
2027    }
2028}
2029
2030fn validate_transcription_temperature(
2031    model: &str,
2032    temperature: Option<f32>,
2033) -> Result<Option<f32>> {
2034    let temperature = validate_temperature(temperature)?;
2035    if model == kcode_openai_api::GPT_4O_TRANSCRIBE && temperature.is_some() {
2036        return Err(Error::invalid(
2037            "OpenAI transcription does not accept a temperature",
2038        ));
2039    }
2040    Ok(temperature)
2041}
2042
2043fn audio_generation_options(
2044    max_output_tokens: u32,
2045    temperature: Option<f32>,
2046) -> Result<GenerationOptions> {
2047    Ok(GenerationOptions {
2048        max_output_tokens: Some(max_output_tokens),
2049        temperature: validate_temperature(temperature)?,
2050        thinking_level: Some(ThinkingLevel::High),
2051        service_tier: ServiceTier::Standard,
2052    })
2053}
2054
2055fn gemini_model(model: &str) -> Option<TextModel> {
2056    match model {
2057        kcode_gemini_api::GEMINI_25_FLASH => Some(TextModel::Flash25),
2058        kcode_gemini_api::GEMINI_31_FLASH_LITE => Some(TextModel::FlashLite),
2059        kcode_gemini_api::GEMINI_31_PRO => Some(TextModel::Pro),
2060        _ => None,
2061    }
2062}
2063
2064fn codex_search_profile(
2065    model: &str,
2066) -> Result<(
2067    kcode_codex_runtime::ReasoningEffort,
2068    kcode_codex_runtime::WebSearchContext,
2069    kcode_codex_runtime::SearchDepth,
2070    Duration,
2071)> {
2072    match model {
2073        QUALITY_SEARCH_MODEL => Ok((
2074            QUALITY_SEARCH_REASONING,
2075            QUALITY_SEARCH_CONTEXT,
2076            QUALITY_SEARCH_DEPTH,
2077            QUALITY_SEARCH_TIMEOUT,
2078        )),
2079        BALANCED_SEARCH_MODEL => Ok((
2080            BALANCED_SEARCH_REASONING,
2081            BALANCED_SEARCH_CONTEXT,
2082            BALANCED_SEARCH_DEPTH,
2083            BALANCED_SEARCH_TIMEOUT,
2084        )),
2085        _ => Err(Error::invalid(
2086            "unsupported exact web-search model; use a supported Gemini model, gpt-5.6-sol, or gpt-5.6-terra",
2087        )),
2088    }
2089}
2090
2091fn normalized_content_type(value: &str) -> String {
2092    value
2093        .split(';')
2094        .next()
2095        .unwrap_or("application/octet-stream")
2096        .trim()
2097        .to_ascii_lowercase()
2098}
2099
2100fn is_ogg(file_name: &str, content_type: &str) -> bool {
2101    matches!(content_type, "audio/ogg" | "video/ogg" | "application/ogg")
2102        || file_name.rsplit_once('.').is_some_and(|(_, extension)| {
2103            matches!(
2104                extension.to_ascii_lowercase().as_str(),
2105                "ogg" | "oga" | "opus"
2106            )
2107        })
2108}
2109
2110fn safe_audio_filename(value: &str, content_type: &str) -> String {
2111    let extension = match content_type {
2112        "audio/ogg" | "audio/opus" | "application/ogg" | "video/ogg" => "ogg",
2113        "audio/wav" | "audio/x-wav" => "wav",
2114        "audio/mpeg" | "audio/mp3" => "mp3",
2115        "audio/mp4" => "mp4",
2116        "audio/webm" => "webm",
2117        "audio/flac" | "audio/x-flac" => "flac",
2118        "audio/m4a" => "m4a",
2119        _ => "audio",
2120    };
2121    let cleaned = value
2122        .chars()
2123        .filter(|character| {
2124            character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_')
2125        })
2126        .take(120)
2127        .collect::<String>();
2128    let supported = cleaned.rsplit_once('.').is_some_and(|(_, extension)| {
2129        matches!(
2130            extension.to_ascii_lowercase().as_str(),
2131            "flac"
2132                | "mp3"
2133                | "mp4"
2134                | "mpeg"
2135                | "mpga"
2136                | "m4a"
2137                | "ogg"
2138                | "oga"
2139                | "opus"
2140                | "wav"
2141                | "webm"
2142        )
2143    });
2144    if cleaned.is_empty() || !supported {
2145        format!("voice-note.{extension}")
2146    } else {
2147        cleaned
2148    }
2149}
2150
2151fn gemini_media(media: &Media) -> Result<GeminiMediaInput> {
2152    match media.kind {
2153        MediaKind::Image => GeminiMediaInput::image(&media.content_type, media.bytes.clone()),
2154        MediaKind::Audio => GeminiMediaInput::audio(&media.content_type, media.bytes.clone()),
2155        MediaKind::Video => GeminiMediaInput::video(&media.content_type, media.bytes.clone()),
2156    }
2157    .map_err(gemini_error)
2158}
2159
2160fn resolved_api_model(
2161    requested: &str,
2162    provider_model: String,
2163    provider: AgentProvider,
2164    context_window_tokens: Option<u64>,
2165    max_input_tokens: Option<u64>,
2166) -> ResolvedAgentModel {
2167    let known = match provider {
2168        AgentProvider::OpenAi if requested == "gpt-5.6" => Some((1_000_000, 700_000)),
2169        AgentProvider::Gemini
2170            if matches!(
2171                requested,
2172                "gemini-2.5-flash" | "gemini-3.1-flash-lite" | "gemini-3.1-pro-preview"
2173            ) =>
2174        {
2175            Some((1_000_000, 700_000))
2176        }
2177        _ => None,
2178    };
2179    let context_window_tokens = context_window_tokens
2180        .or_else(|| known.map(|limits| limits.0))
2181        .unwrap_or(128_000);
2182    let max_input_tokens = max_input_tokens
2183        .or_else(|| known.map(|limits| limits.1))
2184        .unwrap_or_else(|| context_window_tokens.saturating_mul(70) / 100)
2185        .min(context_window_tokens);
2186    ResolvedAgentModel {
2187        requested_model: requested.to_owned(),
2188        provider_model,
2189        provider,
2190        context_window_tokens,
2191        max_input_tokens,
2192    }
2193}
2194
2195fn openai_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> OpenAiAgentRequest {
2196    OpenAiAgentRequest {
2197        model: request.model.clone(),
2198        input: request.input.clone(),
2199        reasoning_effort: request.reasoning_effort.as_str().into(),
2200        tools: request
2201            .tools
2202            .iter()
2203            .map(|tool| OpenAiAgentTool {
2204                name: tool.name.clone(),
2205                description: tool.description.clone(),
2206                input_schema: tool.input_schema.clone(),
2207            })
2208            .collect(),
2209    }
2210}
2211
2212fn gemini_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> GeminiAgentRequest {
2213    GeminiAgentRequest {
2214        model: request.model.clone(),
2215        input: request.input.clone(),
2216        reasoning_effort: request.reasoning_effort.as_str().into(),
2217        tools: request
2218            .tools
2219            .iter()
2220            .map(|tool| GeminiAgentTool {
2221                name: tool.name.clone(),
2222                description: tool.description.clone(),
2223                input_schema: tool.input_schema.clone(),
2224            })
2225            .collect(),
2226    }
2227}
2228
2229fn provider_request_json(request: &kcode_codex_runtime_v2::AgentRequest) -> String {
2230    format!(
2231        "{}\n",
2232        serde_json::to_string(&json!({
2233            "model": request.model,
2234            "input": request.input,
2235            "reasoningEffort": request.reasoning_effort.as_str(),
2236            "tools": request.tools.iter().map(|tool| json!({
2237                "name": tool.name,
2238                "description": tool.description,
2239                "parameters": tool.input_schema
2240            })).collect::<Vec<_>>()
2241        }))
2242        .expect("provider request values always serialize")
2243    )
2244}
2245
2246fn agent_model_context(
2247    request: &kcode_codex_runtime_v2::AgentRequest,
2248    provider: &str,
2249) -> kcode_codex_runtime_v2::ModelContext {
2250    kcode_codex_runtime_v2::ModelContext {
2251        input: request.input.clone(),
2252        provider: provider.into(),
2253        model: request.model.clone(),
2254        reasoning_effort: request.reasoning_effort.as_str().into(),
2255        base_instructions: None,
2256        developer_instructions: None,
2257        tools: request.tools.clone(),
2258    }
2259}
2260
2261fn buffered_openai_turn(
2262    provider_input: String,
2263    model_context: kcode_codex_runtime_v2::ModelContext,
2264    response: kcode_openai_api::AgentTurnResponse,
2265) -> BufferedAgentTurn {
2266    let usage = response
2267        .usage
2268        .as_ref()
2269        .map(|usage| kcode_codex_runtime_v2::TokenUsage {
2270            input_tokens: usage.input_tokens,
2271            output_tokens: usage.output_tokens,
2272            cached_input_tokens: usage.cached_input_tokens,
2273            reasoning_output_tokens: usage.reasoning_output_tokens,
2274            last_input_tokens: None,
2275            last_output_tokens: None,
2276        });
2277    let completed = kcode_codex_runtime_v2::CompletedTurn {
2278        thread_id: response.response_id,
2279        turn_id: Uuid::new_v4().to_string(),
2280        answer: response.text,
2281        usage,
2282    };
2283    buffered_turn(
2284        provider_input,
2285        model_context,
2286        response
2287            .tool_call
2288            .map(|call| kcode_codex_runtime_v2::DynamicToolCall {
2289                call_id: call.call_id,
2290                tool: call.name,
2291                arguments: call.arguments,
2292            }),
2293        completed,
2294    )
2295}
2296
2297fn buffered_gemini_turn(
2298    provider_input: String,
2299    model_context: kcode_codex_runtime_v2::ModelContext,
2300    response: kcode_gemini_api::AgentTurnResponse,
2301) -> BufferedAgentTurn {
2302    let usage = kcode_codex_runtime_v2::TokenUsage {
2303        input_tokens: response.usage.input_tokens,
2304        output_tokens: response
2305            .usage
2306            .output_tokens
2307            .saturating_add(response.usage.thought_tokens),
2308        cached_input_tokens: response.usage.cached_tokens,
2309        reasoning_output_tokens: response.usage.thought_tokens,
2310        last_input_tokens: None,
2311        last_output_tokens: None,
2312    };
2313    let completed = kcode_codex_runtime_v2::CompletedTurn {
2314        thread_id: response.interaction_id,
2315        turn_id: Uuid::new_v4().to_string(),
2316        answer: response.text,
2317        usage: Some(usage),
2318    };
2319    buffered_turn(
2320        provider_input,
2321        model_context,
2322        response
2323            .tool_call
2324            .map(|call| kcode_codex_runtime_v2::DynamicToolCall {
2325                call_id: call.call_id,
2326                tool: call.name,
2327                arguments: call.arguments,
2328            }),
2329        completed,
2330    )
2331}
2332
2333fn buffered_turn(
2334    provider_input: String,
2335    model_context: kcode_codex_runtime_v2::ModelContext,
2336    tool_call: Option<kcode_codex_runtime_v2::DynamicToolCall>,
2337    completed: kcode_codex_runtime_v2::CompletedTurn,
2338) -> BufferedAgentTurn {
2339    let mut events = VecDeque::from([kcode_codex_runtime_v2::AgentEvent::ProviderInput(
2340        provider_input,
2341    )]);
2342    events.push_back(kcode_codex_runtime_v2::AgentEvent::ModelContextSubmitted(
2343        model_context,
2344    ));
2345    if let Some(usage) = completed.usage.clone() {
2346        events.push_back(kcode_codex_runtime_v2::AgentEvent::UsageUpdated(usage));
2347    }
2348    let pending_call_id = tool_call.as_ref().map(|call| call.call_id.clone());
2349    if let Some(call) = tool_call {
2350        events.push_back(kcode_codex_runtime_v2::AgentEvent::ToolCall(call));
2351    } else {
2352        events.push_back(kcode_codex_runtime_v2::AgentEvent::Completed(
2353            completed.clone(),
2354        ));
2355    }
2356    BufferedAgentTurn {
2357        events,
2358        pending_call_id,
2359        completed: Some(completed),
2360    }
2361}
2362
2363fn openai_image_media_type(content_type: &str) -> Result<OpenAiImageMediaType> {
2364    match content_type {
2365        "image/png" => Ok(OpenAiImageMediaType::Png),
2366        "image/jpeg" | "image/jpg" => Ok(OpenAiImageMediaType::Jpeg),
2367        "image/webp" => Ok(OpenAiImageMediaType::WebP),
2368        "image/gif" => Ok(OpenAiImageMediaType::Gif),
2369        _ => Err(Error::invalid(
2370            "OpenAI annotations require PNG, JPEG, WebP, or GIF",
2371        )),
2372    }
2373}
2374
2375fn codex_image_media_type(content_type: &str) -> Result<kcode_codex_runtime_v2::ImageMediaType> {
2376    match content_type {
2377        "image/png" => Ok(kcode_codex_runtime_v2::ImageMediaType::Png),
2378        "image/jpeg" | "image/jpg" => Ok(kcode_codex_runtime_v2::ImageMediaType::Jpeg),
2379        "image/webp" => Ok(kcode_codex_runtime_v2::ImageMediaType::Webp),
2380        _ => Err(Error::invalid(
2381            "Codex annotations require PNG, JPEG, or WebP",
2382        )),
2383    }
2384}
2385
2386fn gemini_usage(usage: &GeminiTokenUsage) -> TokenUsage {
2387    TokenUsage {
2388        input_tokens: usage.input_tokens.saturating_sub(usage.cached_tokens),
2389        cached_input_tokens: usage.cached_tokens,
2390        thinking_tokens: usage.thought_tokens,
2391        output_tokens: usage.output_tokens,
2392    }
2393}
2394
2395fn gemini_cost(cost: &kcode_gemini_api::CostBreakdown) -> CostEstimate {
2396    let accuracy = match cost.accuracy {
2397        kcode_gemini_api::CostAccuracy::Exact => CostAccuracy::Exact,
2398        kcode_gemini_api::CostAccuracy::Estimated => CostAccuracy::Estimated,
2399        kcode_gemini_api::CostAccuracy::Conservative => CostAccuracy::Conservative,
2400    };
2401    CostEstimate {
2402        usd_nanos: cost.total.usd_nanos(),
2403        accuracy,
2404        pricing_version: cost.pricing_version.clone(),
2405    }
2406}
2407
2408fn openai_image_analysis_cost(usage: &ImageAnalysisUsage) -> CostEstimate {
2409    let cached = usage.cached_input_tokens.unwrap_or(0);
2410    let cache_write = usage.cache_write_input_tokens.unwrap_or(0);
2411    let non_cached = usage
2412        .input_tokens
2413        .saturating_sub(cached)
2414        .saturating_sub(cache_write);
2415    let long = usage.input_tokens > 272_000;
2416    let (input_rate, cached_rate, cache_write_rate, output_rate) = if long {
2417        (10_000, 1_000, 12_500, 45_000)
2418    } else {
2419        (5_000, 500, 6_250, 30_000)
2420    };
2421    CostEstimate {
2422        usd_nanos: non_cached
2423            .saturating_mul(input_rate)
2424            .saturating_add(cached.saturating_mul(cached_rate))
2425            .saturating_add(cache_write.saturating_mul(cache_write_rate))
2426            .saturating_add(usage.output_tokens.saturating_mul(output_rate)),
2427        accuracy: CostAccuracy::Exact,
2428        pricing_version: PRICING_VERSION.into(),
2429    }
2430}
2431
2432fn openai_generation_cost(usage: &OpenAiGenerationUsage) -> CostEstimate {
2433    let detailed_input = usage
2434        .input_details
2435        .text_tokens
2436        .saturating_add(usage.input_details.image_tokens);
2437    let output_details = usage.output_details.as_ref();
2438    let detailed_output = output_details
2439        .map(|details| details.text_tokens.saturating_add(details.image_tokens))
2440        .unwrap_or(0);
2441    let input = usage
2442        .input_details
2443        .text_tokens
2444        .saturating_mul(5_000)
2445        .saturating_add(usage.input_details.image_tokens.saturating_mul(8_000))
2446        .saturating_add(
2447            usage
2448                .input_tokens
2449                .saturating_sub(detailed_input)
2450                .saturating_mul(8_000),
2451        );
2452    let output = detailed_output
2453        .saturating_add(usage.output_tokens.saturating_sub(detailed_output))
2454        .saturating_mul(30_000);
2455    CostEstimate {
2456        usd_nanos: input.saturating_add(output),
2457        accuracy: if detailed_input == usage.input_tokens
2458            && output_details.is_some()
2459            && detailed_output == usage.output_tokens
2460        {
2461            CostAccuracy::Exact
2462        } else {
2463            CostAccuracy::Estimated
2464        },
2465        pricing_version: PRICING_VERSION.into(),
2466    }
2467}
2468
2469fn codex_search_cost(model: &str, usage: Option<TokenUsage>) -> Option<CostEstimate> {
2470    let metering = Metering::Tokens(usage?);
2471    estimate_cost(model, &metering).map(|cost| {
2472        // The Codex search wrapper performs at least one OpenAI web-search tool call.
2473        // Provider output does not expose the exact internal search-call count.
2474        cost.with_surcharge(10_000_000, CostAccuracy::Estimated)
2475    })
2476}
2477
2478fn codex_usage(usage: &CodexTokenUsage) -> TokenUsage {
2479    TokenUsage {
2480        input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
2481        cached_input_tokens: usage.cached_input_tokens,
2482        thinking_tokens: usage.reasoning_output_tokens,
2483        output_tokens: usage
2484            .output_tokens
2485            .saturating_sub(usage.reasoning_output_tokens),
2486    }
2487}
2488
2489fn codex_v2_usage(usage: &kcode_codex_runtime_v2::TokenUsage) -> TokenUsage {
2490    TokenUsage {
2491        input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
2492        cached_input_tokens: usage.cached_input_tokens,
2493        thinking_tokens: usage.reasoning_output_tokens,
2494        output_tokens: usage
2495            .output_tokens
2496            .saturating_sub(usage.reasoning_output_tokens),
2497    }
2498}
2499
2500fn openai_image_usage(usage: &ImageAnalysisUsage) -> TokenUsage {
2501    let cached = usage.cached_input_tokens.unwrap_or(0);
2502    let thinking = usage.reasoning_output_tokens.unwrap_or(0);
2503    TokenUsage {
2504        input_tokens: usage.input_tokens.saturating_sub(cached),
2505        cached_input_tokens: cached,
2506        thinking_tokens: thinking,
2507        output_tokens: usage.output_tokens.saturating_sub(thinking),
2508    }
2509}
2510
2511fn openai_generation_usage(usage: &OpenAiGenerationUsage) -> TokenUsage {
2512    TokenUsage {
2513        input_tokens: usage.input_tokens,
2514        cached_input_tokens: 0,
2515        thinking_tokens: 0,
2516        output_tokens: usage.output_tokens,
2517    }
2518}
2519
2520fn transcription_usage(usage: TranscriptionUsage) -> Metering {
2521    match usage {
2522        TranscriptionUsage::DurationSeconds(seconds) => Metering::DurationSeconds { seconds },
2523        TranscriptionUsage::Tokens(tokens) => Metering::Tokens(TokenUsage {
2524            input_tokens: tokens.input_tokens,
2525            cached_input_tokens: 0,
2526            thinking_tokens: 0,
2527            output_tokens: tokens.output_tokens,
2528        }),
2529    }
2530}
2531
2532fn codex_error(error: kcode_codex_runtime::Error) -> Error {
2533    match error.kind() {
2534        CodexErrorKind::InvalidInput => Error::invalid(error.message()),
2535        CodexErrorKind::Authentication => {
2536            Error::unavailable("provider_not_configured", error.message())
2537        }
2538        CodexErrorKind::Unavailable => Error::unavailable("provider_unavailable", error.message()),
2539        CodexErrorKind::RateLimited | CodexErrorKind::Capacity => {
2540            Error::unavailable("provider_rate_limited", error.message())
2541        }
2542        CodexErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
2543        CodexErrorKind::InputTooLarge => Error::invalid(error.message()),
2544        CodexErrorKind::EmptyOutput | CodexErrorKind::Protocol => {
2545            Error::provider("provider_error", error.message())
2546        }
2547    }
2548}
2549
2550fn codex_v2_error(error: kcode_codex_runtime_v2::Error) -> Error {
2551    use kcode_codex_runtime_v2::ErrorKind;
2552    match error.kind() {
2553        ErrorKind::InvalidInput => Error::invalid(error.message()),
2554        ErrorKind::Unavailable | ErrorKind::Authentication => {
2555            Error::unavailable("provider_unavailable", error.message())
2556        }
2557        ErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
2558        ErrorKind::Protocol => Error::provider("provider_error", error.message()),
2559        ErrorKind::Cancelled => Error::cancelled(),
2560    }
2561}
2562
2563fn gemini_error(error: GeminiError) -> Error {
2564    match &error {
2565        GeminiError::InvalidApiKey => {
2566            Error::unavailable("provider_not_configured", error.to_string())
2567        }
2568        GeminiError::InvalidInput(_) => Error::invalid(error.to_string()),
2569        GeminiError::SpendingLimitReached { .. } => {
2570            Error::unavailable("provider_rate_limited", error.to_string())
2571        }
2572        GeminiError::Accounting(_)
2573        | GeminiError::Transport(_)
2574        | GeminiError::Provider { .. }
2575        | GeminiError::Protocol(_) => Error::provider("provider_error", error.to_string()),
2576    }
2577}
2578
2579fn openai_error(error: OpenAiError) -> Error {
2580    match &error {
2581        OpenAiError::InvalidApiKey => {
2582            Error::unavailable("provider_not_configured", error.to_string())
2583        }
2584        OpenAiError::InvalidInput(_) => Error::invalid(error.to_string()),
2585        OpenAiError::Transport(_) | OpenAiError::Provider { .. } | OpenAiError::Protocol(_) => {
2586            Error::provider("provider_error", error.to_string())
2587        }
2588    }
2589}
2590
2591fn web_fetch_error(error: kcode_web_fetch::Error) -> Error {
2592    match error.kind() {
2593        WebFetchErrorKind::InvalidInput | WebFetchErrorKind::UnsafeDestination => {
2594            Error::invalid(error.message())
2595        }
2596        WebFetchErrorKind::Timeout => Error::provider("web_fetch_timeout", error.message()),
2597        WebFetchErrorKind::UnsupportedContent => {
2598            Error::invalid(format!("unsupported web content: {}", error.message()))
2599        }
2600        WebFetchErrorKind::Transport
2601        | WebFetchErrorKind::HttpStatus
2602        | WebFetchErrorKind::EmptyContent => Error::provider("web_fetch_failed", error.message()),
2603    }
2604}
2605
2606fn document_error(error: kcode_doc_extraction::Error) -> Error {
2607    match error.kind() {
2608        DocumentErrorKind::InvalidInput | DocumentErrorKind::UnsupportedFormat => {
2609            Error::invalid(error.message())
2610        }
2611        DocumentErrorKind::ExtractionFailed | DocumentErrorKind::EmptyText => {
2612            Error::provider("document_extraction_failed", error.message())
2613        }
2614    }
2615}
2616
2617#[cfg(test)]
2618mod tests {
2619    use super::*;
2620
2621    #[test]
2622    fn audio_kind_overrides_mislabeled_ogg_video_mime() {
2623        let media = Media::audio(vec![1], "voice.ogg", "video/ogg").unwrap();
2624        assert_eq!(media.kind, MediaKind::Audio);
2625        assert_eq!(media.content_type, "audio/ogg");
2626    }
2627
2628    #[test]
2629    fn exact_search_models_replace_modes() {
2630        assert!(codex_search_profile("gpt-5.6-sol").is_ok());
2631        assert!(codex_search_profile("gpt-5.6-terra").is_ok());
2632        assert!(codex_search_profile("fast").is_err());
2633    }
2634
2635    #[test]
2636    fn long_prompt_is_accepted_and_blank_prompt_is_rejected() {
2637        assert!(validate_prompt(&"x".repeat(1_000_001)).is_ok());
2638        assert!(validate_prompt("  \n\t").is_err());
2639    }
2640
2641    #[test]
2642    fn omitted_temperature_preserves_provider_default() {
2643        let options = audio_generation_options(1_024, None).unwrap();
2644        assert_eq!(options.temperature, None);
2645    }
2646
2647    #[test]
2648    fn exact_zero_temperature_is_preserved() {
2649        let options = audio_generation_options(1_024, Some(0.0)).unwrap();
2650        assert_eq!(options.temperature, Some(0.0));
2651    }
2652
2653    #[test]
2654    fn invalid_and_nonfinite_temperatures_are_rejected_without_receipts() {
2655        for temperature in [Some(-0.1), Some(2.1), Some(f32::NAN), Some(f32::INFINITY)] {
2656            let error = validate_temperature(temperature).err().unwrap();
2657            assert!(error.receipt().is_none());
2658        }
2659    }
2660
2661    #[test]
2662    fn openai_transcription_temperature_is_rejected_without_a_receipt() {
2663        let error =
2664            validate_transcription_temperature(kcode_openai_api::GPT_4O_TRANSCRIBE, Some(0.0))
2665                .err()
2666                .unwrap();
2667        assert!(error.receipt().is_none());
2668    }
2669}