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