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