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