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