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