1#![forbid(unsafe_code)]
8
9mod defaults;
10mod error;
11mod receipts;
12
13use std::{
14 collections::{HashMap, VecDeque},
15 path::PathBuf,
16 sync::{Arc, Mutex},
17 time::{Duration, Instant},
18};
19
20use anyhow::Context;
21use chrono::{DateTime, NaiveDate, Utc};
22pub use kcode_codex_runtime::ReasoningEffort;
23use kcode_codex_runtime::{
24 CatalogCache, Codex, CodexConfig, ErrorKind as CodexErrorKind,
25 GenerationRequest as CodexGenerationRequest, TokenUsage as CodexTokenUsage,
26 WebSearchRequest as CodexSearchRequest,
27};
28use kcode_doc_extraction::{DocumentExtractor, DocumentInput, ErrorKind as DocumentErrorKind};
29use kcode_gemini_api::{
30 AgentTool as GeminiAgentTool, AgentTurnRequest as GeminiAgentRequest,
31 CompletionStatus as GeminiCompletionStatus, Error as GeminiError, Gemini, GenerationOptions,
32 GroundedSearchRequest, MediaInput as GeminiMediaInput, MultimodalRequest, NanoBananaProRequest,
33 ServiceTier, StructuredOutput, TextModel, ThinkingLevel, TokenUsage as GeminiTokenUsage,
34};
35use kcode_openai_api::{
36 AgentTool as OpenAiAgentTool, AgentTurnRequest as OpenAiAgentRequest, AudioInput,
37 Error as OpenAiError, ImageAnalysisRequest, ImageAnalysisStatus as OpenAiImageStatus,
38 ImageAnalysisUsage, ImageEditRequest as OpenAiImageEditRequest,
39 ImageGenerationRequest as OpenAiImageRequest, ImageInput as OpenAiImageInput,
40 ImageMediaType as OpenAiImageMediaType, ImageUsage as OpenAiGenerationUsage, OpenAi,
41 TranscriptionRequest as OpenAiTranscriptionRequest, TranscriptionUsage,
42};
43use kcode_web_fetch::{ErrorKind as WebFetchErrorKind, WebFetcher};
44use serde::{Deserialize, Serialize};
45use serde_json::{Value, json};
46use tokio::sync::watch;
47use uuid::Uuid;
48
49use defaults::*;
50pub use error::{Error, ErrorKind, Result};
51use receipts::ReceiptStore;
52pub use receipts::{DailyUsage, DailyUsageKey, Metering, TokenUsage, UsageReceipt};
53
54pub struct Config {
56 pub openai_api_key: Option<String>,
57 pub gemini_api_key: Option<String>,
58 pub codex_catalog_cache: CatalogCache,
59 pub receipt_directory: PathBuf,
60}
61
62#[derive(Clone, Debug, Eq, PartialEq)]
64pub struct RuntimeModel {
65 pub model: String,
66 pub reasoning_effort: String,
67 pub context_window_tokens: u64,
68 pub max_input_tokens: u64,
69}
70
71#[derive(Clone, Copy, Debug, Eq, PartialEq)]
73pub enum AgentProvider {
74 Codex,
76 OpenAi,
78 Gemini,
80}
81
82#[derive(Clone, Debug, Eq, PartialEq)]
84pub struct ResolvedAgentModel {
85 pub requested_model: String,
87 pub provider_model: String,
89 pub provider: AgentProvider,
91 pub context_window_tokens: u64,
93 pub max_input_tokens: u64,
95}
96
97#[derive(Clone)]
98pub struct Intelligence {
99 codex: Codex,
100 agent: kcode_codex_runtime_v2::Codex,
101 openai: Option<OpenAi>,
102 gemini: Option<Gemini>,
103 web_fetcher: WebFetcher,
104 document_extractor: DocumentExtractor,
105 active_operations: ActiveOperations,
106 receipts: ReceiptStore,
107 model_cache: Arc<Mutex<HashMap<String, ResolvedAgentModel>>>,
108}
109
110#[derive(Clone)]
112pub struct UserIntelligence {
113 service: Intelligence,
114 user_id: String,
115}
116
117#[derive(Clone, Default)]
118struct ActiveOperations {
119 senders: Arc<Mutex<HashMap<Uuid, watch::Sender<bool>>>>,
120}
121
122struct ActiveOperation {
123 id: Uuid,
124 operations: ActiveOperations,
125 cancellation: watch::Receiver<bool>,
126 parent_cancellation: Option<watch::Receiver<bool>>,
127}
128
129enum AgentTurnBackend {
130 Codex(kcode_codex_runtime_v2::AgentTurn),
131 Buffered(BufferedAgentTurn),
132}
133
134struct BufferedAgentTurn {
135 events: VecDeque<kcode_codex_runtime_v2::AgentEvent>,
136 pending_call_id: Option<String>,
137 completed: Option<kcode_codex_runtime_v2::CompletedTurn>,
138}
139
140pub struct AgentTurn {
141 inner: AgentTurnBackend,
142 operation: ActiveOperation,
143 user_id: String,
144 requested_model: String,
145 actual_model: String,
146 provider_request_id: Option<String>,
147 receipts: ReceiptStore,
148 receipt_recorded: bool,
149}
150
151#[derive(Clone, Debug, Eq, PartialEq)]
152pub struct SearchRequest {
153 pub question: String,
154 pub model: String,
155 pub operation_id: Uuid,
156 pub parent_operation_id: Option<Uuid>,
157}
158
159#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
160#[serde(rename_all = "camelCase")]
161pub struct WebSource {
162 pub title: String,
163 pub url: String,
164}
165
166#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
167#[serde(rename_all = "camelCase")]
168pub struct SearchResponse {
169 pub answer: String,
170 pub sources: Vec<WebSource>,
171 pub model: String,
172 pub usage: Option<TokenUsage>,
173}
174
175#[derive(Clone, Debug, Eq, PartialEq)]
176pub struct FetchRequest {
177 pub url: String,
178 pub operation_id: Uuid,
179 pub parent_operation_id: Option<Uuid>,
180}
181
182#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
183#[serde(rename_all = "camelCase")]
184pub struct FetchResponse {
185 pub url: String,
186 pub title: Option<String>,
187 pub content_type: String,
188 pub content: String,
189 pub truncated: bool,
190 pub retrieved_at: DateTime<Utc>,
191}
192
193#[derive(Clone, Copy, Debug, Eq, PartialEq)]
194pub enum MediaKind {
195 Image,
196 Audio,
197 Video,
198}
199
200#[derive(Clone, Debug, Eq, PartialEq)]
201pub struct Media {
202 pub kind: MediaKind,
203 pub bytes: Vec<u8>,
204 pub file_name: String,
205 pub content_type: String,
206}
207
208impl Media {
209 pub fn new(
210 kind: MediaKind,
211 bytes: Vec<u8>,
212 file_name: impl Into<String>,
213 content_type: impl Into<String>,
214 ) -> Result<Self> {
215 if bytes.is_empty() || bytes.len() > MAX_MEDIA_ANNOTATION_BYTES {
216 return Err(Error::invalid(format!(
217 "media must contain between 1 and {MAX_MEDIA_ANNOTATION_BYTES} bytes"
218 )));
219 }
220 let file_name = file_name.into();
221 let mut content_type = normalized_content_type(&content_type.into());
222 if kind == MediaKind::Audio && is_ogg(&file_name, &content_type) {
223 content_type = "audio/ogg".into();
224 }
225 Ok(Self {
226 kind,
227 bytes,
228 file_name,
229 content_type,
230 })
231 }
232
233 pub fn audio(
234 bytes: Vec<u8>,
235 file_name: impl Into<String>,
236 content_type: impl Into<String>,
237 ) -> Result<Self> {
238 Self::new(MediaKind::Audio, bytes, file_name, content_type)
239 }
240}
241
242#[derive(Clone, Debug, Eq, PartialEq)]
243pub struct TranscriptionRequest {
244 pub prompt: String,
245 pub model: String,
246 pub media: Media,
247 pub operation_id: Uuid,
248 pub parent_operation_id: Option<Uuid>,
249}
250
251#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
252#[serde(rename_all = "camelCase")]
253pub struct TranscriptionResponse {
254 pub model: String,
255 pub text: String,
256 pub metering: Metering,
257}
258
259#[derive(Clone, Debug, PartialEq)]
261pub struct StructuredAudioRequest {
262 pub operation: String,
263 pub prompt: String,
264 pub model: String,
265 pub media: Media,
266 pub schema: Value,
267 pub max_output_tokens: u32,
268 pub operation_id: Uuid,
269 pub parent_operation_id: Option<Uuid>,
270}
271
272#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
273#[serde(rename_all = "camelCase")]
274pub struct StructuredAudioResponse {
275 pub model: String,
276 pub text: String,
277 pub usage: TokenUsage,
278}
279
280#[derive(Clone, Debug, Eq, PartialEq)]
282pub struct TextGenerationRequest {
283 pub operation: String,
284 pub prompt: String,
285 pub model: String,
286 pub reasoning_effort: ReasoningEffort,
287 pub timeout: Duration,
288 pub operation_id: Uuid,
289 pub parent_operation_id: Option<Uuid>,
290}
291
292#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
293#[serde(rename_all = "camelCase")]
294pub struct TextGenerationResponse {
295 pub model: String,
296 pub text: String,
297 pub thread_id: String,
298 pub usage: Option<TokenUsage>,
299}
300
301#[derive(Clone, Debug, Eq, PartialEq)]
302pub struct AnnotationRequest {
303 pub prompt: String,
304 pub model: String,
305 pub media: Media,
306 pub operation_id: Uuid,
307 pub parent_operation_id: Option<Uuid>,
308}
309
310#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
311#[serde(rename_all = "camelCase")]
312pub struct AnnotationResponse {
313 pub complete: bool,
314 pub model: String,
315 pub file_name: String,
316 pub content_type: String,
317 pub text: String,
318 pub incomplete_reason: Option<String>,
319 pub usage: Option<TokenUsage>,
320}
321
322#[derive(Clone, Debug, Eq, PartialEq)]
324pub struct ImageRequest {
325 pub model: String,
327 pub prompt: String,
329 pub references: Vec<Media>,
331 pub operation_id: Uuid,
333 pub parent_operation_id: Option<Uuid>,
335}
336
337#[derive(Clone, Debug, Eq, PartialEq)]
339pub struct ImageResponse {
340 pub model: String,
342 pub content_type: String,
344 pub bytes: Vec<u8>,
346 pub usage: Option<TokenUsage>,
348}
349
350#[derive(Clone, Debug, Eq, PartialEq)]
351pub struct Document {
352 pub bytes: Vec<u8>,
353 pub file_name: String,
354 pub content_type: String,
355}
356
357#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
358#[serde(rename_all = "camelCase")]
359pub struct DocumentExtraction {
360 pub file_name: String,
361 pub content_type: String,
362 pub format: String,
363 pub text: String,
364 pub characters: usize,
365 pub truncated: bool,
366}
367
368pub async fn open(config: Config) -> anyhow::Result<(Intelligence, RuntimeModel)> {
369 let mut codex_config = CodexConfig::new(DEFAULT_MODEL);
370 codex_config.base_instruction = KENNEDY_CODEX_BASE_INSTRUCTION.into();
371 codex_config.validation_reasoning_effort = GENERATION_REASONING_EFFORT;
372 let codex = Codex::open(codex_config, config.codex_catalog_cache)
373 .await
374 .context("opening Kennedy Codex runtime")?;
375 let limits = codex
376 .catalog()
377 .model_limits(DEFAULT_MODEL)
378 .with_context(|| format!("Codex model {DEFAULT_MODEL} is absent from the catalog"))?;
379 let runtime = RuntimeModel {
380 model: DEFAULT_MODEL.into(),
381 reasoning_effort: GENERATION_REASONING_EFFORT.as_str().into(),
382 context_window_tokens: limits.context_window_tokens(),
383 max_input_tokens: limits.max_input_tokens(),
384 };
385 let agent_config = kcode_codex_runtime_v2::CodexConfig {
386 executable: codex.catalog().executable().to_owned(),
387 base_instruction: KENNEDY_CODEX_BASE_INSTRUCTION.into(),
388 model_catalog: Some(codex.catalog().path().to_owned()),
389 ..kcode_codex_runtime_v2::CodexConfig::default()
390 };
391 let agent = kcode_codex_runtime_v2::Codex::open(agent_config)
392 .await
393 .context("opening Kennedy Codex runtime v2")?;
394 let openai = config
395 .openai_api_key
396 .filter(|value| !value.trim().is_empty())
397 .map(OpenAi::open)
398 .transpose()
399 .context("opening OpenAI client")?;
400 let gemini = config
401 .gemini_api_key
402 .filter(|value| !value.trim().is_empty())
403 .map(Gemini::open)
404 .transpose()
405 .context("opening Gemini client")?;
406 Ok((
407 Intelligence {
408 codex,
409 agent,
410 openai,
411 gemini,
412 web_fetcher: WebFetcher::default(),
413 document_extractor: DocumentExtractor::default(),
414 active_operations: ActiveOperations::default(),
415 receipts: ReceiptStore::open(config.receipt_directory)
416 .map_err(anyhow::Error::new)
417 .context("opening intelligence usage receipts")?,
418 model_cache: Arc::new(Mutex::new(HashMap::new())),
419 },
420 runtime,
421 ))
422}
423
424impl Intelligence {
425 pub fn for_user(&self, user_id: impl Into<String>) -> Result<UserIntelligence> {
426 let user_id = user_id.into();
427 if user_id.trim().is_empty() || user_id.chars().count() > 256 {
428 return Err(Error::invalid(
429 "user_id must contain between 1 and 256 characters",
430 ));
431 }
432 Ok(UserIntelligence {
433 service: self.clone(),
434 user_id,
435 })
436 }
437
438 pub fn cancel(&self, operation_id: Uuid) -> Result<bool> {
439 self.active_operations.cancel(operation_id)
440 }
441
442 pub fn receipts(&self) -> Result<Vec<UsageReceipt>> {
443 self.receipts.receipts()
444 }
445
446 pub fn daily_usage(
447 &self,
448 day: NaiveDate,
449 ) -> Result<std::collections::BTreeMap<DailyUsageKey, DailyUsage>> {
450 self.receipts.daily_usage(day)
451 }
452
453 pub async fn extract_document(&self, document: Document) -> Result<DocumentExtraction> {
454 extract_document(&self.document_extractor, document).await
455 }
456
457 pub async fn resolve_agent_model(&self, requested: &str) -> Result<ResolvedAgentModel> {
459 validate_agent_model(requested)?;
460 if let Some(resolved) = self
461 .model_cache
462 .lock()
463 .map_err(|_| {
464 Error::internal("model_cache_unavailable", "The model cache is unavailable.")
465 })?
466 .get(requested)
467 .cloned()
468 {
469 return Ok(resolved);
470 }
471 let resolved = if let Some(provider_model) = requested.strip_prefix("codex/") {
472 let limits = self
473 .codex
474 .catalog()
475 .model_limits(provider_model)
476 .ok_or_else(|| {
477 Error::invalid(format!(
478 "{provider_model:?} is not an available model in the Codex catalog"
479 ))
480 })?;
481 ResolvedAgentModel {
482 requested_model: requested.to_owned(),
483 provider_model: provider_model.to_owned(),
484 provider: AgentProvider::Codex,
485 context_window_tokens: limits.context_window_tokens(),
486 max_input_tokens: limits.max_input_tokens(),
487 }
488 } else if requested.ends_with("-sol")
489 || requested.ends_with("-terra")
490 || requested.ends_with("-luna")
491 {
492 let limits = self
493 .codex
494 .catalog()
495 .model_limits(requested)
496 .ok_or_else(|| {
497 Error::invalid(format!(
498 "{requested:?} is not an available model in the Codex catalog"
499 ))
500 })?;
501 ResolvedAgentModel {
502 requested_model: requested.to_owned(),
503 provider_model: requested.to_owned(),
504 provider: AgentProvider::Codex,
505 context_window_tokens: limits.context_window_tokens(),
506 max_input_tokens: limits.max_input_tokens(),
507 }
508 } else if requested.starts_with("gemini-") {
509 let gemini = self.gemini.as_ref().ok_or_else(|| {
510 Error::unavailable("provider_not_configured", "Gemini is not configured.")
511 })?;
512 let metadata =
513 tokio::time::timeout(MODEL_DISCOVERY_TIMEOUT, gemini.model_metadata(requested))
514 .await
515 .map_err(|_| {
516 Error::provider("provider_timeout", "Gemini model discovery timed out.")
517 })?
518 .map_err(gemini_error)?;
519 resolved_api_model(
520 requested,
521 metadata.id,
522 AgentProvider::Gemini,
523 metadata.context_window_tokens,
524 metadata.max_input_tokens,
525 )
526 } else {
527 let openai = self.openai.as_ref().ok_or_else(|| {
528 Error::unavailable("provider_not_configured", "OpenAI is not configured.")
529 })?;
530 let metadata =
531 tokio::time::timeout(MODEL_DISCOVERY_TIMEOUT, openai.model_metadata(requested))
532 .await
533 .map_err(|_| {
534 Error::provider("provider_timeout", "OpenAI model discovery timed out.")
535 })?
536 .map_err(openai_error)?;
537 resolved_api_model(
538 requested,
539 metadata.id,
540 AgentProvider::OpenAi,
541 metadata.context_window_tokens,
542 metadata.max_input_tokens,
543 )
544 };
545 self.model_cache
546 .lock()
547 .map_err(|_| {
548 Error::internal("model_cache_unavailable", "The model cache is unavailable.")
549 })?
550 .insert(requested.to_owned(), resolved.clone());
551 Ok(resolved)
552 }
553}
554
555impl UserIntelligence {
556 pub fn user_id(&self) -> &str {
557 &self.user_id
558 }
559
560 pub async fn start_agent_turn(
561 &self,
562 operation_id: Uuid,
563 parent_operation_id: Option<Uuid>,
564 mut request: kcode_codex_runtime_v2::AgentRequest,
565 ) -> Result<AgentTurn> {
566 let resolved = self.service.resolve_agent_model(&request.model).await?;
567 if request.previous_thread_id.is_some() {
568 return Err(Error::invalid(
569 "agent turns must submit complete input without a previous thread; this keeps each usage receipt scoped to one call",
570 ));
571 }
572 let mut operation = self
573 .service
574 .active_operations
575 .register_request(operation_id, parent_operation_id)?;
576 let requested_model = resolved.requested_model.clone();
577 request.model = resolved.provider_model.clone();
578 let (inner, actual_model, provider_request_id) = match resolved.provider {
579 AgentProvider::Codex => {
580 let inner = self.account_result(
581 "agent_turn",
582 &requested_model,
583 self.service
584 .agent
585 .start_turn(request)
586 .await
587 .map_err(codex_v2_error),
588 )?;
589 (
590 AgentTurnBackend::Codex(inner),
591 resolved.provider_model,
592 None,
593 )
594 }
595 AgentProvider::OpenAi => {
596 let openai = self.service.openai.as_ref().ok_or_else(|| {
597 Error::unavailable("provider_not_configured", "OpenAI is not configured.")
598 })?;
599 let provider_request = openai_agent_request(&request);
600 let provider_input = provider_request_json(&request);
601 let result = tokio::select! {
602 _ = operation.cancelled() => Err(Error::cancelled()),
603 result = tokio::time::timeout(request.timeout, openai.agent_turn(provider_request)) => {
604 result
605 .map_err(|_| Error::provider("provider_timeout", "OpenAI agent turn timed out."))
606 .and_then(|result| result.map_err(openai_error))
607 }
608 };
609 let result = self.account_result("agent_turn", &requested_model, result)?;
610 let actual_model = result.model.clone();
611 let provider_request_id = Some(result.response_id.clone());
612 (
613 AgentTurnBackend::Buffered(buffered_openai_turn(provider_input, result)),
614 actual_model,
615 provider_request_id,
616 )
617 }
618 AgentProvider::Gemini => {
619 let gemini = self.service.gemini.as_ref().ok_or_else(|| {
620 Error::unavailable("provider_not_configured", "Gemini is not configured.")
621 })?;
622 let provider_request = gemini_agent_request(&request);
623 let provider_input = provider_request_json(&request);
624 let result = tokio::select! {
625 _ = operation.cancelled() => Err(Error::cancelled()),
626 result = tokio::time::timeout(request.timeout, gemini.agent_turn(provider_request)) => {
627 result
628 .map_err(|_| Error::provider("provider_timeout", "Gemini agent turn timed out."))
629 .and_then(|result| result.map_err(gemini_error))
630 }
631 };
632 let result = self.account_result("agent_turn", &requested_model, result)?;
633 let actual_model = result.model.clone();
634 let provider_request_id = Some(result.interaction_id.clone());
635 (
636 AgentTurnBackend::Buffered(buffered_gemini_turn(provider_input, result)),
637 actual_model,
638 provider_request_id,
639 )
640 }
641 };
642 Ok(AgentTurn {
643 inner,
644 operation,
645 user_id: self.user_id.clone(),
646 requested_model,
647 actual_model,
648 provider_request_id,
649 receipts: self.service.receipts.clone(),
650 receipt_recorded: false,
651 })
652 }
653
654 pub async fn search(&self, request: SearchRequest) -> Result<SearchResponse> {
655 let question = request.question.trim();
656 if question.is_empty() || question.chars().count() > 4_000 {
657 return Err(Error::invalid(
658 "question must contain between 1 and 4000 characters",
659 ));
660 }
661 validate_model(&request.model)?;
662 let mut operation = self
663 .service
664 .active_operations
665 .register_request(request.operation_id, request.parent_operation_id)?;
666 let started = Instant::now();
667 let response = if let Some(model) = gemini_model(&request.model) {
668 let gemini = self.service.gemini.as_ref().ok_or_else(|| {
669 Error::unavailable(
670 "provider_not_configured",
671 "Gemini search is not configured.",
672 )
673 })?;
674 let result = tokio::select! {
675 _ = operation.cancelled() => Err(Error::cancelled()),
676 result = tokio::time::timeout(
677 FAST_SEARCH_TIMEOUT,
678 gemini.grounded_search_with_model(
679 model,
680 GroundedSearchRequest::new(question),
681 ),
682 ) => result
683 .map_err(|_| Error::provider("provider_timeout", "Gemini search timed out."))
684 .and_then(|result| result.map_err(gemini_error)),
685 };
686 let result = self.account_result("web_search", &request.model, result)?;
687 let interaction = result.interaction;
688 let usage = gemini_usage(&interaction.usage);
689 self.record_tokens(
690 "web_search",
691 &request.model,
692 &interaction.model,
693 Some(usage),
694 Some(interaction.id.clone()),
695 None,
696 )?;
697 SearchResponse {
698 answer: interaction
699 .text
700 .filter(|value| !value.trim().is_empty())
701 .ok_or_else(|| {
702 Error::provider("provider_error", "Gemini search returned no answer text.")
703 })?,
704 sources: result
705 .sources
706 .into_iter()
707 .map(|source| WebSource {
708 title: source.title,
709 url: source.url,
710 })
711 .collect(),
712 model: interaction.model,
713 usage: Some(usage),
714 }
715 } else {
716 let (reasoning, context, depth, timeout) = codex_search_profile(&request.model)?;
717 let result = tokio::select! {
718 _ = operation.cancelled() => Err(Error::cancelled()),
719 result = self.service.codex.web_search(CodexSearchRequest {
720 question: question.to_owned(),
721 model: request.model.clone(),
722 reasoning_effort: reasoning,
723 context,
724 depth,
725 timeout,
726 }) => result.map_err(codex_error),
727 };
728 let result = self.account_result("web_search", &request.model, result)?;
729 let usage = result.usage.as_ref().map(codex_usage);
730 self.record_tokens(
731 "web_search",
732 &request.model,
733 &request.model,
734 usage,
735 None,
736 None,
737 )?;
738 SearchResponse {
739 answer: result.answer,
740 sources: result
741 .sources
742 .into_iter()
743 .map(|source| WebSource {
744 title: source.title,
745 url: source.url,
746 })
747 .collect(),
748 model: request.model.clone(),
749 usage,
750 }
751 };
752 tracing::info!(
753 user_id = %self.user_id,
754 operation_id = %request.operation_id,
755 model = %response.model,
756 duration_ms = started.elapsed().as_millis(),
757 "Intelligence search completed"
758 );
759 Ok(response)
760 }
761
762 pub async fn transcribe_structured_audio(
764 &self,
765 request: StructuredAudioRequest,
766 ) -> Result<StructuredAudioResponse> {
767 validate_operation(&request.operation)?;
768 validate_prompt(&request.prompt)?;
769 validate_model(&request.model)?;
770 if request.media.kind != MediaKind::Audio {
771 return Err(Error::invalid(
772 "structured audio inference requires audio media",
773 ));
774 }
775 if request.max_output_tokens == 0 || request.max_output_tokens > 65_536 {
776 return Err(Error::invalid(
777 "max_output_tokens must be between 1 and 65536",
778 ));
779 }
780 let model = gemini_model(&request.model).ok_or_else(|| {
781 Error::invalid("structured audio inference requires a supported exact Gemini model")
782 })?;
783 let gemini = self.service.gemini.as_ref().ok_or_else(|| {
784 Error::unavailable(
785 "provider_not_configured",
786 "Gemini structured audio inference is not configured.",
787 )
788 })?;
789 let media = gemini_media(&request.media)?;
790 let structured_output = StructuredOutput::new(request.schema).map_err(gemini_error)?;
791 let mut provider_request = MultimodalRequest::new(request.prompt, vec![media]);
792 provider_request.options = GenerationOptions {
793 max_output_tokens: Some(request.max_output_tokens),
794 temperature: None,
795 thinking_level: Some(ThinkingLevel::High),
796 service_tier: ServiceTier::Standard,
797 };
798 provider_request.structured_output = Some(structured_output);
799 let mut operation = self
800 .service
801 .active_operations
802 .register_request(request.operation_id, request.parent_operation_id)?;
803 let result = tokio::select! {
804 _ = operation.cancelled() => Err(Error::cancelled()),
805 result = tokio::time::timeout(
806 MEDIA_ANNOTATION_TIMEOUT,
807 gemini.infer_multimodal(model, provider_request),
808 ) => result
809 .map_err(|_| Error::provider("provider_timeout", "Gemini structured audio inference timed out."))
810 .and_then(|result| result.map_err(gemini_error)),
811 };
812 let result = self.account_result(&request.operation, &request.model, result)?;
813 let usage = gemini_usage(&result.usage);
814 self.record_tokens(
815 &request.operation,
816 &request.model,
817 &result.model,
818 Some(usage),
819 Some(result.id.clone()),
820 None,
821 )?;
822 if result.status != GeminiCompletionStatus::Completed {
823 return Err(Error::provider(
824 "provider_incomplete",
825 "Gemini structured audio inference did not complete.",
826 ));
827 }
828 let text = result
829 .text
830 .filter(|text| !text.trim().is_empty())
831 .ok_or_else(|| {
832 Error::provider(
833 "provider_empty_output",
834 "Gemini structured audio inference returned no text.",
835 )
836 })?;
837 Ok(StructuredAudioResponse {
838 model: result.model,
839 text,
840 usage,
841 })
842 }
843
844 pub async fn generate_text(
846 &self,
847 request: TextGenerationRequest,
848 ) -> Result<TextGenerationResponse> {
849 validate_operation(&request.operation)?;
850 validate_model(&request.model)?;
851 if request.prompt.trim().is_empty() || request.prompt.chars().count() > 1_000_000 {
852 return Err(Error::invalid(
853 "generation prompt must contain between 1 and 1000000 characters",
854 ));
855 }
856 if request.timeout.is_zero() || request.timeout > Duration::from_secs(60 * 60) {
857 return Err(Error::invalid(
858 "generation timeout must be between 1 second and 1 hour",
859 ));
860 }
861 let mut operation = self
862 .service
863 .active_operations
864 .register_request(request.operation_id, request.parent_operation_id)?;
865 let mut provider_request =
866 CodexGenerationRequest::new(request.prompt, request.model.clone());
867 provider_request.reasoning_effort = request.reasoning_effort;
868 provider_request.ephemeral = true;
869 provider_request.timeout = request.timeout;
870 let result = tokio::select! {
871 _ = operation.cancelled() => Err(Error::cancelled()),
872 result = self.service.codex.generate(provider_request) => result.map_err(codex_error),
873 };
874 let result = self.account_result(&request.operation, &request.model, result)?;
875 let usage = result.usage.as_ref().map(codex_usage);
876 self.record_tokens(
877 &request.operation,
878 &request.model,
879 &request.model,
880 usage,
881 None,
882 Some(result.thread_id.clone()),
883 )?;
884 Ok(TextGenerationResponse {
885 model: request.model,
886 text: result.answer,
887 thread_id: result.thread_id,
888 usage,
889 })
890 }
891
892 pub async fn fetch(&self, request: FetchRequest) -> Result<FetchResponse> {
893 let mut operation = self
894 .service
895 .active_operations
896 .register_request(request.operation_id, request.parent_operation_id)?;
897 let fetched = tokio::select! {
898 _ = operation.cancelled() => Err(Error::cancelled()),
899 result = self.service.web_fetcher.fetch(&request.url) => result.map_err(web_fetch_error),
900 }?;
901 Ok(FetchResponse {
902 url: fetched.url,
903 title: fetched.title,
904 content_type: fetched.content_type,
905 content: fetched.content,
906 truncated: fetched.truncated,
907 retrieved_at: DateTime::<Utc>::from(fetched.retrieved_at),
908 })
909 }
910
911 pub async fn transcribe(&self, request: TranscriptionRequest) -> Result<TranscriptionResponse> {
912 validate_prompt(&request.prompt)?;
913 if request.media.kind != MediaKind::Audio {
914 return Err(Error::invalid("transcription requires audio media"));
915 }
916 validate_model(&request.model)?;
917 let mut operation = self
918 .service
919 .active_operations
920 .register_request(request.operation_id, request.parent_operation_id)?;
921 let response = if request.model == kcode_openai_api::GPT_4O_TRANSCRIBE {
922 let openai = self.service.openai.as_ref().ok_or_else(|| {
923 Error::unavailable(
924 "transcription_unavailable",
925 "OpenAI audio transcription is not configured.",
926 )
927 })?;
928 let input = AudioInput::new(
929 safe_audio_filename(&request.media.file_name, &request.media.content_type),
930 request.media.content_type.clone(),
931 request.media.bytes,
932 )
933 .map_err(openai_error)?;
934 let mut provider_request = OpenAiTranscriptionRequest::new(input);
935 provider_request.prompt = Some(request.prompt);
936 let result = tokio::select! {
937 _ = operation.cancelled() => Err(Error::cancelled()),
938 result = tokio::time::timeout(
939 MEDIA_ANNOTATION_TIMEOUT,
940 openai.transcribe(provider_request),
941 ) => result
942 .map_err(|_| Error::provider("provider_timeout", "OpenAI transcription timed out."))
943 .and_then(|result| result.map_err(openai_error)),
944 };
945 let result = self.account_result("transcribe_audio", &request.model, result)?;
946 let metering = result
947 .usage
948 .map(transcription_usage)
949 .unwrap_or(Metering::Unavailable);
950 self.record_metering(
951 "transcribe_audio",
952 &request.model,
953 &request.model,
954 metering.clone(),
955 None,
956 None,
957 )?;
958 TranscriptionResponse {
959 model: request.model.clone(),
960 text: result.text,
961 metering,
962 }
963 } else {
964 let model = gemini_model(&request.model).ok_or_else(|| {
965 Error::invalid("transcription model must be gpt-4o-transcribe or a supported exact Gemini model")
966 })?;
967 let gemini = self.service.gemini.as_ref().ok_or_else(|| {
968 Error::unavailable(
969 "provider_not_configured",
970 "Gemini transcription is not configured.",
971 )
972 })?;
973 let media = gemini_media(&request.media)?;
974 let result = tokio::select! {
975 _ = operation.cancelled() => Err(Error::cancelled()),
976 result = tokio::time::timeout(
977 MEDIA_ANNOTATION_TIMEOUT,
978 gemini.infer_multimodal(
979 model,
980 MultimodalRequest::new(request.prompt, vec![media]),
981 ),
982 ) => result
983 .map_err(|_| Error::provider("provider_timeout", "Gemini transcription timed out."))
984 .and_then(|result| result.map_err(gemini_error)),
985 };
986 let result = self.account_result("transcribe_audio", &request.model, result)?;
987 let metering = Metering::Tokens(gemini_usage(&result.usage));
988 self.record_metering(
989 "transcribe_audio",
990 &request.model,
991 &result.model,
992 metering.clone(),
993 Some(result.id.clone()),
994 None,
995 )?;
996 TranscriptionResponse {
997 model: result.model,
998 text: result
999 .text
1000 .filter(|text| !text.trim().is_empty())
1001 .ok_or_else(|| {
1002 Error::provider(
1003 "empty_transcription",
1004 "Gemini returned no transcription text.",
1005 )
1006 })?,
1007 metering,
1008 }
1009 };
1010 Ok(response)
1011 }
1012
1013 pub async fn annotate(&self, request: AnnotationRequest) -> Result<AnnotationResponse> {
1014 validate_prompt(&request.prompt)?;
1015 validate_model(&request.model)?;
1016 let mut operation = self
1017 .service
1018 .active_operations
1019 .register_request(request.operation_id, request.parent_operation_id)?;
1020 let file_name = request.media.file_name.clone();
1021 let content_type = request.media.content_type.clone();
1022 let response = if let Some(model) = gemini_model(&request.model) {
1023 let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1024 Error::unavailable(
1025 "provider_not_configured",
1026 "Gemini media annotation is not configured.",
1027 )
1028 })?;
1029 let media = gemini_media(&request.media)?;
1030 let result = tokio::select! {
1031 _ = operation.cancelled() => Err(Error::cancelled()),
1032 result = tokio::time::timeout(
1033 MEDIA_ANNOTATION_TIMEOUT,
1034 gemini.infer_multimodal(
1035 model,
1036 MultimodalRequest::new(request.prompt, vec![media]),
1037 ),
1038 ) => result
1039 .map_err(|_| Error::provider("provider_timeout", "Gemini annotation timed out."))
1040 .and_then(|result| result.map_err(gemini_error)),
1041 };
1042 let result = self.account_result("annotate_media", &request.model, result)?;
1043 let usage = gemini_usage(&result.usage);
1044 self.record_tokens(
1045 "annotate_media",
1046 &request.model,
1047 &result.model,
1048 Some(usage),
1049 Some(result.id.clone()),
1050 None,
1051 )?;
1052 AnnotationResponse {
1053 complete: result.status == GeminiCompletionStatus::Completed,
1054 model: result.model,
1055 file_name,
1056 content_type,
1057 text: result
1058 .text
1059 .filter(|text| !text.trim().is_empty())
1060 .ok_or_else(|| {
1061 Error::provider("empty_annotation", "Gemini returned no annotation text.")
1062 })?,
1063 incomplete_reason: None,
1064 usage: Some(usage),
1065 }
1066 } else if request.model == "gpt-5.6" {
1067 if request.media.kind != MediaKind::Image {
1068 return Err(Error::invalid(
1069 "OpenAI media annotation accepts images only",
1070 ));
1071 }
1072 let openai = self.service.openai.as_ref().ok_or_else(|| {
1073 Error::unavailable(
1074 "provider_not_configured",
1075 "OpenAI media annotation is not configured.",
1076 )
1077 })?;
1078 let image = OpenAiImageInput::new(
1079 openai_image_media_type(&request.media.content_type)?,
1080 request.media.bytes,
1081 )
1082 .map_err(openai_error)?;
1083 let result = tokio::select! {
1084 _ = operation.cancelled() => Err(Error::cancelled()),
1085 result = tokio::time::timeout(
1086 MEDIA_ANNOTATION_TIMEOUT,
1087 openai.analyze_image(ImageAnalysisRequest::new(image, request.prompt)),
1088 ) => result
1089 .map_err(|_| Error::provider("provider_timeout", "OpenAI annotation timed out."))
1090 .and_then(|result| result.map_err(openai_error)),
1091 };
1092 let result = self.account_result("annotate_media", &request.model, result)?;
1093 let usage = result.usage.as_ref().map(openai_image_usage);
1094 self.record_tokens(
1095 "annotate_media",
1096 &request.model,
1097 &result.model,
1098 usage,
1099 None,
1100 None,
1101 )?;
1102 let (complete, incomplete_reason) = match result.status {
1103 OpenAiImageStatus::Completed => (true, None),
1104 OpenAiImageStatus::Incomplete { reason } => (false, reason),
1105 };
1106 AnnotationResponse {
1107 complete,
1108 model: result.model,
1109 file_name,
1110 content_type,
1111 text: result.text,
1112 incomplete_reason,
1113 usage,
1114 }
1115 } else if matches!(
1116 request.model.as_str(),
1117 "gpt-5.6-sol" | "gpt-5.6-terra" | "gpt-5.6-luna"
1118 ) {
1119 if request.media.kind != MediaKind::Image {
1120 return Err(Error::invalid("Codex media annotation accepts images only"));
1121 }
1122 let image = kcode_codex_runtime_v2::ImageInput::new(
1123 codex_image_media_type(&request.media.content_type)?,
1124 request.media.bytes,
1125 )
1126 .map_err(codex_v2_error)?;
1127 let turn_result = self
1128 .service
1129 .agent
1130 .start_image_turn(kcode_codex_runtime_v2::ImageTurnRequest::new(
1131 request.prompt,
1132 request.model.clone(),
1133 vec![image],
1134 ))
1135 .await
1136 .map_err(codex_v2_error);
1137 let mut turn = self.account_result("annotate_media", &request.model, turn_result)?;
1138 let completed = loop {
1139 let event = tokio::select! {
1140 _ = operation.cancelled() => {
1141 turn.cancel();
1142 self.record_metering(
1143 "annotate_media",
1144 &request.model,
1145 &request.model,
1146 Metering::Unavailable,
1147 None,
1148 None,
1149 )?;
1150 return Err(Error::cancelled());
1151 }
1152 event = turn.next_event() => event,
1153 };
1154 match event {
1155 Some(Ok(kcode_codex_runtime_v2::AgentEvent::ProviderInput(_))) => {}
1156 Some(Ok(kcode_codex_runtime_v2::AgentEvent::ToolCall(_))) => {
1157 turn.cancel();
1158 self.record_metering(
1159 "annotate_media",
1160 &request.model,
1161 &request.model,
1162 Metering::Unavailable,
1163 None,
1164 None,
1165 )?;
1166 return Err(Error::provider(
1167 "provider_error",
1168 "A tool-free Codex image turn requested a tool.",
1169 ));
1170 }
1171 Some(Ok(kcode_codex_runtime_v2::AgentEvent::Completed(completed))) => {
1172 break completed;
1173 }
1174 Some(Err(error)) => {
1175 self.record_metering(
1176 "annotate_media",
1177 &request.model,
1178 &request.model,
1179 Metering::Unavailable,
1180 None,
1181 None,
1182 )?;
1183 return Err(codex_v2_error(error));
1184 }
1185 None => {
1186 self.record_metering(
1187 "annotate_media",
1188 &request.model,
1189 &request.model,
1190 Metering::Unavailable,
1191 None,
1192 None,
1193 )?;
1194 return Err(Error::provider(
1195 "empty_annotation",
1196 "Codex ended without annotation text.",
1197 ));
1198 }
1199 }
1200 };
1201 let usage = completed.usage.as_ref().map(codex_v2_usage);
1202 self.record_tokens(
1203 "annotate_media",
1204 &request.model,
1205 &request.model,
1206 usage,
1207 Some(completed.turn_id.clone()),
1208 Some(completed.thread_id.clone()),
1209 )?;
1210 AnnotationResponse {
1211 complete: true,
1212 model: request.model.clone(),
1213 file_name,
1214 content_type,
1215 text: completed.answer,
1216 incomplete_reason: None,
1217 usage,
1218 }
1219 } else {
1220 return Err(Error::invalid(format!(
1221 "unsupported exact annotation model {}",
1222 request.model
1223 )));
1224 };
1225 Ok(response)
1226 }
1227
1228 pub async fn generate_image(&self, request: ImageRequest) -> Result<ImageResponse> {
1230 validate_model(&request.model)?;
1231 if request.prompt.trim().is_empty()
1232 || request.prompt.chars().count() > MAX_IMAGE_PROMPT_CHARACTERS
1233 {
1234 return Err(Error::invalid(format!(
1235 "image prompt must contain 1 through {MAX_IMAGE_PROMPT_CHARACTERS} characters"
1236 )));
1237 }
1238 if request
1239 .references
1240 .iter()
1241 .any(|media| media.kind != MediaKind::Image)
1242 {
1243 return Err(Error::invalid("image references must all be images"));
1244 }
1245 let mut operation = self
1246 .service
1247 .active_operations
1248 .register_request(request.operation_id, request.parent_operation_id)?;
1249 let operation_name = if request.references.is_empty() {
1250 "generate_image"
1251 } else {
1252 "edit_image"
1253 };
1254 if request.model == kcode_openai_api::GPT_IMAGE_2 {
1255 let openai = self.service.openai.as_ref().ok_or_else(|| {
1256 Error::unavailable(
1257 "provider_not_configured",
1258 "OpenAI image generation is not configured.",
1259 )
1260 })?;
1261 let result = if request.references.is_empty() {
1262 let provider_request = OpenAiImageRequest::new(request.prompt);
1263 tokio::select! {
1264 _ = operation.cancelled() => Err(Error::cancelled()),
1265 result = tokio::time::timeout(
1266 IMAGE_OPERATION_TIMEOUT,
1267 openai.generate_image(provider_request),
1268 ) => result
1269 .map_err(|_| Error::provider("provider_timeout", "OpenAI image generation timed out."))
1270 .and_then(|result| result.map_err(openai_error)),
1271 }
1272 } else {
1273 let mut images = request
1274 .references
1275 .into_iter()
1276 .map(|media| {
1277 OpenAiImageInput::new(
1278 openai_image_media_type(&media.content_type)?,
1279 media.bytes,
1280 )
1281 .map_err(openai_error)
1282 })
1283 .collect::<Result<Vec<_>>>()?;
1284 let first = images.remove(0);
1285 let mut provider_request = OpenAiImageEditRequest::new(first, request.prompt);
1286 provider_request.images.extend(images);
1287 tokio::select! {
1288 _ = operation.cancelled() => Err(Error::cancelled()),
1289 result = tokio::time::timeout(
1290 IMAGE_OPERATION_TIMEOUT,
1291 openai.edit_image(provider_request),
1292 ) => result
1293 .map_err(|_| Error::provider("provider_timeout", "OpenAI image editing timed out."))
1294 .and_then(|result| result.map_err(openai_error)),
1295 }
1296 };
1297 let result = self.account_result(operation_name, &request.model, result)?;
1298 let usage = result.usage.as_ref().map(openai_generation_usage);
1299 self.record_tokens(
1300 operation_name,
1301 &request.model,
1302 kcode_openai_api::GPT_IMAGE_2,
1303 usage,
1304 result.request_id,
1305 None,
1306 )?;
1307 Ok(ImageResponse {
1308 model: kcode_openai_api::GPT_IMAGE_2.into(),
1309 content_type: result.image.format.mime_type().into(),
1310 bytes: result.image.data,
1311 usage,
1312 })
1313 } else if request.model == kcode_gemini_api::NANO_BANANA_PRO {
1314 let gemini = self.service.gemini.as_ref().ok_or_else(|| {
1315 Error::unavailable(
1316 "provider_not_configured",
1317 "Gemini image generation is not configured.",
1318 )
1319 })?;
1320 let mut provider_request = NanoBananaProRequest::new(request.prompt);
1321 provider_request.images = request
1322 .references
1323 .into_iter()
1324 .map(|media| {
1325 GeminiMediaInput::image(&media.content_type, media.bytes).map_err(gemini_error)
1326 })
1327 .collect::<Result<Vec<_>>>()?;
1328 let result = tokio::select! {
1329 _ = operation.cancelled() => Err(Error::cancelled()),
1330 result = tokio::time::timeout(
1331 IMAGE_OPERATION_TIMEOUT,
1332 gemini.nano_banana_pro(provider_request),
1333 ) => result
1334 .map_err(|_| Error::provider("provider_timeout", "Gemini image generation timed out."))
1335 .and_then(|result| result.map_err(gemini_error)),
1336 };
1337 let result = self.account_result(operation_name, &request.model, result)?;
1338 let usage = gemini_usage(&result.usage);
1339 self.record_tokens(
1340 operation_name,
1341 &request.model,
1342 &result.model,
1343 Some(usage),
1344 Some(result.id.clone()),
1345 None,
1346 )?;
1347 let mut images = result.images;
1348 if images.len() != 1 {
1349 return Err(Error::provider(
1350 "provider_error",
1351 "Gemini did not return exactly one generated image.",
1352 ));
1353 }
1354 let image = images.remove(0);
1355 Ok(ImageResponse {
1356 model: result.model,
1357 content_type: image.mime_type,
1358 bytes: image.data,
1359 usage: Some(usage),
1360 })
1361 } else {
1362 Err(Error::invalid(format!(
1363 "unsupported exact image model {}",
1364 request.model
1365 )))
1366 }
1367 }
1368
1369 fn record_tokens(
1370 &self,
1371 operation: &str,
1372 requested_model: &str,
1373 actual_model: &str,
1374 usage: Option<TokenUsage>,
1375 provider_request_id: Option<String>,
1376 provider_thread_id: Option<String>,
1377 ) -> Result<()> {
1378 self.record_metering(
1379 operation,
1380 requested_model,
1381 actual_model,
1382 usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
1383 provider_request_id,
1384 provider_thread_id,
1385 )
1386 }
1387
1388 fn record_metering(
1389 &self,
1390 operation: &str,
1391 requested_model: &str,
1392 actual_model: &str,
1393 metering: Metering,
1394 provider_request_id: Option<String>,
1395 provider_thread_id: Option<String>,
1396 ) -> Result<()> {
1397 let mut receipt = UsageReceipt::new(
1398 self.user_id.clone(),
1399 operation,
1400 requested_model,
1401 actual_model,
1402 metering,
1403 );
1404 receipt.provider_request_id = provider_request_id;
1405 receipt.provider_thread_id = provider_thread_id;
1406 self.service.receipts.record(&receipt)
1407 }
1408
1409 fn account_result<T>(
1410 &self,
1411 operation: &str,
1412 requested_model: &str,
1413 result: Result<T>,
1414 ) -> Result<T> {
1415 match result {
1416 Ok(value) => Ok(value),
1417 Err(error) => {
1418 self.record_metering(
1419 operation,
1420 requested_model,
1421 requested_model,
1422 Metering::Unavailable,
1423 None,
1424 None,
1425 )?;
1426 Err(error)
1427 }
1428 }
1429 }
1430}
1431
1432async fn extract_document(
1433 extractor: &DocumentExtractor,
1434 document: Document,
1435) -> Result<DocumentExtraction> {
1436 let extractor = extractor.clone();
1437 let extracted = tokio::task::spawn_blocking(move || {
1438 extractor.extract(DocumentInput {
1439 file_name: document.file_name,
1440 content_type: document.content_type,
1441 data: document.bytes,
1442 })
1443 })
1444 .await
1445 .map_err(|_| {
1446 Error::internal(
1447 "document_extraction_failed",
1448 "The document extraction worker stopped unexpectedly.",
1449 )
1450 })?
1451 .map_err(document_error)?;
1452 Ok(DocumentExtraction {
1453 file_name: extracted.file_name,
1454 content_type: extracted.content_type,
1455 format: extracted.format.as_str().into(),
1456 text: extracted.text,
1457 characters: extracted.characters,
1458 truncated: extracted.truncated,
1459 })
1460}
1461
1462impl AgentTurn {
1463 pub async fn next_event(&mut self) -> Result<Option<kcode_codex_runtime_v2::AgentEvent>> {
1464 let event = match &mut self.inner {
1465 AgentTurnBackend::Codex(inner) => tokio::select! {
1466 _ = self.operation.cancelled() => {
1467 inner.cancel();
1468 self.record_unavailable()?;
1469 return Err(Error::cancelled());
1470 }
1471 event = inner.next_event() => match event {
1472 Some(Ok(event)) => Some(event),
1473 Some(Err(error)) => {
1474 self.record_unavailable()?;
1475 return Err(codex_v2_error(error));
1476 }
1477 None => {
1478 self.record_unavailable()?;
1479 None
1480 },
1481 }
1482 },
1483 AgentTurnBackend::Buffered(inner) => {
1484 if *self.operation.cancellation.borrow() {
1485 self.record_unavailable()?;
1486 return Err(Error::cancelled());
1487 }
1488 inner.events.pop_front()
1489 }
1490 };
1491 if let Some(kcode_codex_runtime_v2::AgentEvent::Completed(completed)) = &event
1492 && !self.receipt_recorded
1493 {
1494 let mut receipt = UsageReceipt::new(
1495 self.user_id.clone(),
1496 "agent_turn",
1497 self.requested_model.clone(),
1498 self.actual_model.clone(),
1499 completed
1500 .usage
1501 .as_ref()
1502 .map(codex_v2_usage)
1503 .map(Metering::Tokens)
1504 .unwrap_or(Metering::Unavailable),
1505 );
1506 receipt.provider_request_id = self
1507 .provider_request_id
1508 .clone()
1509 .or_else(|| Some(completed.turn_id.clone()));
1510 receipt.provider_thread_id = Some(completed.thread_id.clone());
1511 self.receipts.record(&receipt)?;
1512 self.receipt_recorded = true;
1513 }
1514 Ok(event)
1515 }
1516
1517 fn record_unavailable(&mut self) -> Result<()> {
1518 if self.receipt_recorded {
1519 return Ok(());
1520 }
1521 let receipt = UsageReceipt::new(
1522 self.user_id.clone(),
1523 "agent_turn",
1524 self.requested_model.clone(),
1525 self.actual_model.clone(),
1526 Metering::Unavailable,
1527 );
1528 self.receipts.record(&receipt)?;
1529 self.receipt_recorded = true;
1530 Ok(())
1531 }
1532
1533 pub async fn respond(
1534 &mut self,
1535 call_id: &str,
1536 result: kcode_codex_runtime_v2::ToolResult,
1537 ) -> Result<()> {
1538 match &mut self.inner {
1539 AgentTurnBackend::Codex(inner) => {
1540 inner.respond(call_id, result).await.map_err(codex_v2_error)
1541 }
1542 AgentTurnBackend::Buffered(inner) => {
1543 if inner.pending_call_id.as_deref() != Some(call_id) {
1544 return Err(Error::invalid(
1545 "tool result does not match the pending provider call",
1546 ));
1547 }
1548 inner.pending_call_id = None;
1549 if let Some(mut completed) = inner.completed.take() {
1550 completed.answer.clear();
1551 inner
1552 .events
1553 .push_back(kcode_codex_runtime_v2::AgentEvent::Completed(completed));
1554 }
1555 Ok(())
1556 }
1557 }
1558 }
1559}
1560
1561impl Drop for AgentTurn {
1562 fn drop(&mut self) {
1563 let _ = self.record_unavailable();
1564 }
1565}
1566
1567impl ActiveOperations {
1568 fn register(&self, id: Uuid) -> Result<ActiveOperation> {
1569 let (sender, cancellation) = watch::channel(false);
1570 let mut senders = self.senders.lock().map_err(|_| {
1571 Error::internal(
1572 "operation_registry_unavailable",
1573 "The operation registry is unavailable.",
1574 )
1575 })?;
1576 if senders.contains_key(&id) {
1577 return Err(Error::conflict(
1578 "operation_in_progress",
1579 "An operation with this identifier is already running.",
1580 ));
1581 }
1582 senders.insert(id, sender);
1583 Ok(ActiveOperation {
1584 id,
1585 operations: self.clone(),
1586 cancellation,
1587 parent_cancellation: None,
1588 })
1589 }
1590
1591 fn register_request(
1592 &self,
1593 request_id: Uuid,
1594 parent_operation_id: Option<Uuid>,
1595 ) -> Result<ActiveOperation> {
1596 if let Some(parent_id) = parent_operation_id {
1597 if request_id == parent_id {
1598 return Err(Error::invalid(
1599 "operation_id and parent_operation_id must be different",
1600 ));
1601 }
1602 let parent_cancellation = self
1603 .senders
1604 .lock()
1605 .map_err(|_| {
1606 Error::internal(
1607 "operation_registry_unavailable",
1608 "The operation registry is unavailable.",
1609 )
1610 })?
1611 .get(&parent_id)
1612 .map(watch::Sender::subscribe)
1613 .ok_or_else(|| {
1614 Error::conflict(
1615 "parent_operation_not_running",
1616 "The parent operation is no longer running.",
1617 )
1618 })?;
1619 let mut operation = self.register(request_id)?;
1620 operation.parent_cancellation = Some(parent_cancellation);
1621 Ok(operation)
1622 } else {
1623 self.register(request_id)
1624 }
1625 }
1626
1627 fn cancel(&self, id: Uuid) -> Result<bool> {
1628 let sender = self
1629 .senders
1630 .lock()
1631 .map_err(|_| {
1632 Error::internal(
1633 "operation_registry_unavailable",
1634 "The operation registry is unavailable.",
1635 )
1636 })?
1637 .get(&id)
1638 .cloned();
1639 Ok(sender.is_some_and(|sender| sender.send(true).is_ok()))
1640 }
1641
1642 fn remove(&self, id: Uuid) {
1643 if let Ok(mut senders) = self.senders.lock() {
1644 senders.remove(&id);
1645 }
1646 }
1647}
1648
1649impl ActiveOperation {
1650 async fn cancelled(&mut self) {
1651 if let Some(parent) = &mut self.parent_cancellation {
1652 tokio::select! {
1653 _ = cancellation_requested(&mut self.cancellation) => {}
1654 _ = cancellation_requested(parent) => {}
1655 }
1656 } else {
1657 cancellation_requested(&mut self.cancellation).await;
1658 }
1659 }
1660}
1661
1662async fn cancellation_requested(cancellation: &mut watch::Receiver<bool>) {
1663 if *cancellation.borrow() {
1664 return;
1665 }
1666 while cancellation.changed().await.is_ok() {
1667 if *cancellation.borrow() {
1668 return;
1669 }
1670 }
1671}
1672
1673impl Drop for ActiveOperation {
1674 fn drop(&mut self) {
1675 self.operations.remove(self.id);
1676 }
1677}
1678
1679fn validate_model(model: &str) -> Result<()> {
1680 if model.trim().is_empty()
1681 || model.chars().count() > 128
1682 || !model
1683 .bytes()
1684 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_'))
1685 {
1686 return Err(Error::invalid(
1687 "model must be an exact safe model identifier",
1688 ));
1689 }
1690 Ok(())
1691}
1692
1693fn validate_agent_model(model: &str) -> Result<()> {
1694 if model.trim().is_empty()
1695 || model.chars().count() > 128
1696 || !model.bytes().all(|byte| {
1697 byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_' | b'/' | b':')
1698 })
1699 || model.starts_with("codex:")
1700 {
1701 return Err(Error::invalid(
1702 "model must be an exact safe provider model identifier",
1703 ));
1704 }
1705 Ok(())
1706}
1707
1708fn validate_operation(operation: &str) -> Result<()> {
1709 if operation.is_empty()
1710 || operation.len() > 64
1711 || !operation
1712 .bytes()
1713 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
1714 {
1715 return Err(Error::invalid(
1716 "operation must be a lowercase identifier of at most 64 bytes",
1717 ));
1718 }
1719 Ok(())
1720}
1721
1722fn validate_prompt(prompt: &str) -> Result<()> {
1723 if prompt.trim().is_empty() || prompt.chars().count() > MAX_MEDIA_ANNOTATION_PROMPT_CHARACTERS {
1724 return Err(Error::invalid(format!(
1725 "prompt must contain between 1 and {MAX_MEDIA_ANNOTATION_PROMPT_CHARACTERS} characters"
1726 )));
1727 }
1728 Ok(())
1729}
1730
1731fn gemini_model(model: &str) -> Option<TextModel> {
1732 match model {
1733 kcode_gemini_api::GEMINI_25_FLASH => Some(TextModel::Flash25),
1734 kcode_gemini_api::GEMINI_31_FLASH_LITE => Some(TextModel::FlashLite),
1735 kcode_gemini_api::GEMINI_31_PRO => Some(TextModel::Pro),
1736 _ => None,
1737 }
1738}
1739
1740fn codex_search_profile(
1741 model: &str,
1742) -> Result<(
1743 kcode_codex_runtime::ReasoningEffort,
1744 kcode_codex_runtime::WebSearchContext,
1745 kcode_codex_runtime::SearchDepth,
1746 Duration,
1747)> {
1748 match model {
1749 QUALITY_SEARCH_MODEL => Ok((
1750 QUALITY_SEARCH_REASONING,
1751 QUALITY_SEARCH_CONTEXT,
1752 QUALITY_SEARCH_DEPTH,
1753 QUALITY_SEARCH_TIMEOUT,
1754 )),
1755 BALANCED_SEARCH_MODEL => Ok((
1756 BALANCED_SEARCH_REASONING,
1757 BALANCED_SEARCH_CONTEXT,
1758 BALANCED_SEARCH_DEPTH,
1759 BALANCED_SEARCH_TIMEOUT,
1760 )),
1761 _ => Err(Error::invalid(
1762 "unsupported exact web-search model; use a supported Gemini model, gpt-5.6-sol, or gpt-5.6-terra",
1763 )),
1764 }
1765}
1766
1767fn normalized_content_type(value: &str) -> String {
1768 value
1769 .split(';')
1770 .next()
1771 .unwrap_or("application/octet-stream")
1772 .trim()
1773 .to_ascii_lowercase()
1774}
1775
1776fn is_ogg(file_name: &str, content_type: &str) -> bool {
1777 matches!(content_type, "audio/ogg" | "video/ogg" | "application/ogg")
1778 || file_name.rsplit_once('.').is_some_and(|(_, extension)| {
1779 matches!(
1780 extension.to_ascii_lowercase().as_str(),
1781 "ogg" | "oga" | "opus"
1782 )
1783 })
1784}
1785
1786fn safe_audio_filename(value: &str, content_type: &str) -> String {
1787 let extension = match content_type {
1788 "audio/ogg" | "audio/opus" | "application/ogg" | "video/ogg" => "ogg",
1789 "audio/wav" | "audio/x-wav" => "wav",
1790 "audio/mpeg" | "audio/mp3" => "mp3",
1791 "audio/mp4" => "mp4",
1792 "audio/webm" => "webm",
1793 "audio/flac" | "audio/x-flac" => "flac",
1794 "audio/m4a" => "m4a",
1795 _ => "audio",
1796 };
1797 let cleaned = value
1798 .chars()
1799 .filter(|character| {
1800 character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_')
1801 })
1802 .take(120)
1803 .collect::<String>();
1804 let supported = cleaned.rsplit_once('.').is_some_and(|(_, extension)| {
1805 matches!(
1806 extension.to_ascii_lowercase().as_str(),
1807 "flac"
1808 | "mp3"
1809 | "mp4"
1810 | "mpeg"
1811 | "mpga"
1812 | "m4a"
1813 | "ogg"
1814 | "oga"
1815 | "opus"
1816 | "wav"
1817 | "webm"
1818 )
1819 });
1820 if cleaned.is_empty() || !supported {
1821 format!("voice-note.{extension}")
1822 } else {
1823 cleaned
1824 }
1825}
1826
1827fn gemini_media(media: &Media) -> Result<GeminiMediaInput> {
1828 match media.kind {
1829 MediaKind::Image => GeminiMediaInput::image(&media.content_type, media.bytes.clone()),
1830 MediaKind::Audio => GeminiMediaInput::audio(&media.content_type, media.bytes.clone()),
1831 MediaKind::Video => GeminiMediaInput::video(&media.content_type, media.bytes.clone()),
1832 }
1833 .map_err(gemini_error)
1834}
1835
1836fn resolved_api_model(
1837 requested: &str,
1838 provider_model: String,
1839 provider: AgentProvider,
1840 context_window_tokens: Option<u64>,
1841 max_input_tokens: Option<u64>,
1842) -> ResolvedAgentModel {
1843 let known = match provider {
1844 AgentProvider::OpenAi if requested == "gpt-5.6" => Some((1_000_000, 700_000)),
1845 AgentProvider::Gemini
1846 if matches!(
1847 requested,
1848 "gemini-2.5-flash" | "gemini-3.1-flash-lite" | "gemini-3.1-pro-preview"
1849 ) =>
1850 {
1851 Some((1_000_000, 700_000))
1852 }
1853 _ => None,
1854 };
1855 let context_window_tokens = context_window_tokens
1856 .or_else(|| known.map(|limits| limits.0))
1857 .unwrap_or(128_000);
1858 let max_input_tokens = max_input_tokens
1859 .or_else(|| known.map(|limits| limits.1))
1860 .unwrap_or_else(|| context_window_tokens.saturating_mul(70) / 100)
1861 .min(context_window_tokens);
1862 ResolvedAgentModel {
1863 requested_model: requested.to_owned(),
1864 provider_model,
1865 provider,
1866 context_window_tokens,
1867 max_input_tokens,
1868 }
1869}
1870
1871fn openai_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> OpenAiAgentRequest {
1872 OpenAiAgentRequest {
1873 model: request.model.clone(),
1874 input: request.input.clone(),
1875 reasoning_effort: request.reasoning_effort.as_str().into(),
1876 tools: request
1877 .tools
1878 .iter()
1879 .map(|tool| OpenAiAgentTool {
1880 name: tool.name.clone(),
1881 description: tool.description.clone(),
1882 input_schema: tool.input_schema.clone(),
1883 })
1884 .collect(),
1885 }
1886}
1887
1888fn gemini_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> GeminiAgentRequest {
1889 GeminiAgentRequest {
1890 model: request.model.clone(),
1891 input: request.input.clone(),
1892 reasoning_effort: request.reasoning_effort.as_str().into(),
1893 tools: request
1894 .tools
1895 .iter()
1896 .map(|tool| GeminiAgentTool {
1897 name: tool.name.clone(),
1898 description: tool.description.clone(),
1899 input_schema: tool.input_schema.clone(),
1900 })
1901 .collect(),
1902 }
1903}
1904
1905fn provider_request_json(request: &kcode_codex_runtime_v2::AgentRequest) -> String {
1906 format!(
1907 "{}\n",
1908 serde_json::to_string(&json!({
1909 "model": request.model,
1910 "input": request.input,
1911 "reasoningEffort": request.reasoning_effort.as_str(),
1912 "tools": request.tools.iter().map(|tool| json!({
1913 "name": tool.name,
1914 "description": tool.description,
1915 "parameters": tool.input_schema
1916 })).collect::<Vec<_>>()
1917 }))
1918 .expect("provider request values always serialize")
1919 )
1920}
1921
1922fn buffered_openai_turn(
1923 provider_input: String,
1924 response: kcode_openai_api::AgentTurnResponse,
1925) -> BufferedAgentTurn {
1926 let usage = response
1927 .usage
1928 .as_ref()
1929 .map(|usage| kcode_codex_runtime_v2::TokenUsage {
1930 input_tokens: usage.input_tokens,
1931 output_tokens: usage.output_tokens,
1932 cached_input_tokens: usage.cached_input_tokens,
1933 reasoning_output_tokens: usage.reasoning_output_tokens,
1934 last_input_tokens: None,
1935 last_output_tokens: None,
1936 });
1937 let completed = kcode_codex_runtime_v2::CompletedTurn {
1938 thread_id: response.response_id,
1939 turn_id: Uuid::new_v4().to_string(),
1940 answer: response.text,
1941 usage,
1942 };
1943 buffered_turn(
1944 provider_input,
1945 response
1946 .tool_call
1947 .map(|call| kcode_codex_runtime_v2::DynamicToolCall {
1948 call_id: call.call_id,
1949 tool: call.name,
1950 arguments: call.arguments,
1951 }),
1952 completed,
1953 )
1954}
1955
1956fn buffered_gemini_turn(
1957 provider_input: String,
1958 response: kcode_gemini_api::AgentTurnResponse,
1959) -> BufferedAgentTurn {
1960 let usage = kcode_codex_runtime_v2::TokenUsage {
1961 input_tokens: response.usage.input_tokens,
1962 output_tokens: response.usage.output_tokens,
1963 cached_input_tokens: response.usage.cached_tokens,
1964 reasoning_output_tokens: response.usage.thought_tokens,
1965 last_input_tokens: None,
1966 last_output_tokens: None,
1967 };
1968 let completed = kcode_codex_runtime_v2::CompletedTurn {
1969 thread_id: response.interaction_id,
1970 turn_id: Uuid::new_v4().to_string(),
1971 answer: response.text,
1972 usage: Some(usage),
1973 };
1974 buffered_turn(
1975 provider_input,
1976 response
1977 .tool_call
1978 .map(|call| kcode_codex_runtime_v2::DynamicToolCall {
1979 call_id: call.call_id,
1980 tool: call.name,
1981 arguments: call.arguments,
1982 }),
1983 completed,
1984 )
1985}
1986
1987fn buffered_turn(
1988 provider_input: String,
1989 tool_call: Option<kcode_codex_runtime_v2::DynamicToolCall>,
1990 completed: kcode_codex_runtime_v2::CompletedTurn,
1991) -> BufferedAgentTurn {
1992 let mut events = VecDeque::from([kcode_codex_runtime_v2::AgentEvent::ProviderInput(
1993 provider_input,
1994 )]);
1995 let pending_call_id = tool_call.as_ref().map(|call| call.call_id.clone());
1996 if let Some(call) = tool_call {
1997 events.push_back(kcode_codex_runtime_v2::AgentEvent::ToolCall(call));
1998 } else {
1999 events.push_back(kcode_codex_runtime_v2::AgentEvent::Completed(
2000 completed.clone(),
2001 ));
2002 }
2003 BufferedAgentTurn {
2004 events,
2005 pending_call_id,
2006 completed: Some(completed),
2007 }
2008}
2009
2010fn openai_image_media_type(content_type: &str) -> Result<OpenAiImageMediaType> {
2011 match content_type {
2012 "image/png" => Ok(OpenAiImageMediaType::Png),
2013 "image/jpeg" | "image/jpg" => Ok(OpenAiImageMediaType::Jpeg),
2014 "image/webp" => Ok(OpenAiImageMediaType::WebP),
2015 "image/gif" => Ok(OpenAiImageMediaType::Gif),
2016 _ => Err(Error::invalid(
2017 "OpenAI annotations require PNG, JPEG, WebP, or GIF",
2018 )),
2019 }
2020}
2021
2022fn codex_image_media_type(content_type: &str) -> Result<kcode_codex_runtime_v2::ImageMediaType> {
2023 match content_type {
2024 "image/png" => Ok(kcode_codex_runtime_v2::ImageMediaType::Png),
2025 "image/jpeg" | "image/jpg" => Ok(kcode_codex_runtime_v2::ImageMediaType::Jpeg),
2026 "image/webp" => Ok(kcode_codex_runtime_v2::ImageMediaType::Webp),
2027 _ => Err(Error::invalid(
2028 "Codex annotations require PNG, JPEG, or WebP",
2029 )),
2030 }
2031}
2032
2033fn gemini_usage(usage: &GeminiTokenUsage) -> TokenUsage {
2034 TokenUsage {
2035 input_tokens: usage.input_tokens.saturating_sub(usage.cached_tokens),
2036 cached_input_tokens: usage.cached_tokens,
2037 thinking_tokens: usage.thought_tokens,
2038 output_tokens: usage.output_tokens,
2039 }
2040}
2041
2042fn codex_usage(usage: &CodexTokenUsage) -> TokenUsage {
2043 TokenUsage {
2044 input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
2045 cached_input_tokens: usage.cached_input_tokens,
2046 thinking_tokens: usage.reasoning_output_tokens,
2047 output_tokens: usage
2048 .output_tokens
2049 .saturating_sub(usage.reasoning_output_tokens),
2050 }
2051}
2052
2053fn codex_v2_usage(usage: &kcode_codex_runtime_v2::TokenUsage) -> TokenUsage {
2054 TokenUsage {
2055 input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
2056 cached_input_tokens: usage.cached_input_tokens,
2057 thinking_tokens: usage.reasoning_output_tokens,
2058 output_tokens: usage
2059 .output_tokens
2060 .saturating_sub(usage.reasoning_output_tokens),
2061 }
2062}
2063
2064fn openai_image_usage(usage: &ImageAnalysisUsage) -> TokenUsage {
2065 let cached = usage.cached_input_tokens.unwrap_or(0);
2066 let thinking = usage.reasoning_output_tokens.unwrap_or(0);
2067 TokenUsage {
2068 input_tokens: usage.input_tokens.saturating_sub(cached),
2069 cached_input_tokens: cached,
2070 thinking_tokens: thinking,
2071 output_tokens: usage.output_tokens.saturating_sub(thinking),
2072 }
2073}
2074
2075fn openai_generation_usage(usage: &OpenAiGenerationUsage) -> TokenUsage {
2076 TokenUsage {
2077 input_tokens: usage.input_tokens,
2078 cached_input_tokens: 0,
2079 thinking_tokens: 0,
2080 output_tokens: usage.output_tokens,
2081 }
2082}
2083
2084fn transcription_usage(usage: TranscriptionUsage) -> Metering {
2085 match usage {
2086 TranscriptionUsage::DurationSeconds(seconds) => Metering::DurationSeconds { seconds },
2087 TranscriptionUsage::Tokens(tokens) => Metering::Tokens(TokenUsage {
2088 input_tokens: tokens.input_tokens,
2089 cached_input_tokens: 0,
2090 thinking_tokens: 0,
2091 output_tokens: tokens.output_tokens,
2092 }),
2093 }
2094}
2095
2096fn codex_error(error: kcode_codex_runtime::Error) -> Error {
2097 match error.kind() {
2098 CodexErrorKind::InvalidInput => Error::invalid(error.message()),
2099 CodexErrorKind::Authentication => {
2100 Error::unavailable("provider_not_configured", error.message())
2101 }
2102 CodexErrorKind::Unavailable => Error::unavailable("provider_unavailable", error.message()),
2103 CodexErrorKind::RateLimited | CodexErrorKind::Capacity => {
2104 Error::unavailable("provider_rate_limited", error.message())
2105 }
2106 CodexErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
2107 CodexErrorKind::InputTooLarge => Error::invalid(error.message()),
2108 CodexErrorKind::EmptyOutput | CodexErrorKind::Protocol => {
2109 Error::provider("provider_error", error.message())
2110 }
2111 }
2112}
2113
2114fn codex_v2_error(error: kcode_codex_runtime_v2::Error) -> Error {
2115 use kcode_codex_runtime_v2::ErrorKind;
2116 match error.kind() {
2117 ErrorKind::InvalidInput => Error::invalid(error.message()),
2118 ErrorKind::Unavailable | ErrorKind::Authentication => {
2119 Error::unavailable("provider_unavailable", error.message())
2120 }
2121 ErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
2122 ErrorKind::Protocol => Error::provider("provider_error", error.message()),
2123 ErrorKind::Cancelled => Error::cancelled(),
2124 }
2125}
2126
2127fn gemini_error(error: GeminiError) -> Error {
2128 match &error {
2129 GeminiError::InvalidApiKey => {
2130 Error::unavailable("provider_not_configured", error.to_string())
2131 }
2132 GeminiError::InvalidInput(_) => Error::invalid(error.to_string()),
2133 GeminiError::SpendingLimitReached { .. } => {
2134 Error::unavailable("provider_rate_limited", error.to_string())
2135 }
2136 GeminiError::Accounting(_)
2137 | GeminiError::Transport(_)
2138 | GeminiError::Provider { .. }
2139 | GeminiError::Protocol(_) => Error::provider("provider_error", error.to_string()),
2140 }
2141}
2142
2143fn openai_error(error: OpenAiError) -> Error {
2144 match &error {
2145 OpenAiError::InvalidApiKey => {
2146 Error::unavailable("provider_not_configured", error.to_string())
2147 }
2148 OpenAiError::InvalidInput(_) => Error::invalid(error.to_string()),
2149 OpenAiError::Transport(_) | OpenAiError::Provider { .. } | OpenAiError::Protocol(_) => {
2150 Error::provider("provider_error", error.to_string())
2151 }
2152 }
2153}
2154
2155fn web_fetch_error(error: kcode_web_fetch::Error) -> Error {
2156 match error.kind() {
2157 WebFetchErrorKind::InvalidInput | WebFetchErrorKind::UnsafeDestination => {
2158 Error::invalid(error.message())
2159 }
2160 WebFetchErrorKind::Timeout => Error::provider("web_fetch_timeout", error.message()),
2161 WebFetchErrorKind::UnsupportedContent => {
2162 Error::invalid(format!("unsupported web content: {}", error.message()))
2163 }
2164 WebFetchErrorKind::Transport
2165 | WebFetchErrorKind::HttpStatus
2166 | WebFetchErrorKind::EmptyContent => Error::provider("web_fetch_failed", error.message()),
2167 }
2168}
2169
2170fn document_error(error: kcode_doc_extraction::Error) -> Error {
2171 match error.kind() {
2172 DocumentErrorKind::InvalidInput | DocumentErrorKind::UnsupportedFormat => {
2173 Error::invalid(error.message())
2174 }
2175 DocumentErrorKind::ExtractionFailed | DocumentErrorKind::EmptyText => {
2176 Error::provider("document_extraction_failed", error.message())
2177 }
2178 }
2179}
2180
2181#[cfg(test)]
2182mod tests {
2183 use super::*;
2184
2185 #[test]
2186 fn audio_kind_overrides_mislabeled_ogg_video_mime() {
2187 let media = Media::audio(vec![1], "voice.ogg", "video/ogg").unwrap();
2188 assert_eq!(media.kind, MediaKind::Audio);
2189 assert_eq!(media.content_type, "audio/ogg");
2190 }
2191
2192 #[test]
2193 fn exact_search_models_replace_modes() {
2194 assert!(codex_search_profile("gpt-5.6-sol").is_ok());
2195 assert!(codex_search_profile("gpt-5.6-terra").is_ok());
2196 assert!(codex_search_profile("fast").is_err());
2197 }
2198
2199 #[tokio::test]
2200 async fn child_operations_have_independent_ids_and_inherit_parent_cancellation() {
2201 let operations = ActiveOperations::default();
2202 let parent_id = Uuid::new_v4();
2203 let child_id = Uuid::new_v4();
2204 let _parent = operations.register(parent_id).unwrap();
2205 let mut child = operations
2206 .register_request(child_id, Some(parent_id))
2207 .unwrap();
2208
2209 assert!(operations.cancel(child_id).unwrap());
2210 tokio::time::timeout(Duration::from_millis(50), child.cancelled())
2211 .await
2212 .unwrap();
2213 assert!(operations.cancel(parent_id).unwrap());
2214
2215 let child_id = Uuid::new_v4();
2216 let mut inherited = operations
2217 .register_request(child_id, Some(parent_id))
2218 .unwrap();
2219 tokio::time::timeout(Duration::from_millis(50), inherited.cancelled())
2220 .await
2221 .unwrap();
2222 }
2223}