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