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