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