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