#![forbid(unsafe_code)]
mod defaults;
mod error;
mod receipts;
use std::{
collections::HashMap,
path::PathBuf,
sync::{Arc, Mutex},
time::{Duration, Instant},
};
use anyhow::Context;
use chrono::{DateTime, NaiveDate, Utc};
pub use kcode_codex_runtime::ReasoningEffort;
use kcode_codex_runtime::{
CatalogCache, Codex, CodexConfig, ErrorKind as CodexErrorKind,
GenerationRequest as CodexGenerationRequest, TokenUsage as CodexTokenUsage,
WebSearchRequest as CodexSearchRequest,
};
use kcode_doc_extraction::{DocumentExtractor, DocumentInput, ErrorKind as DocumentErrorKind};
use kcode_gemini_api::{
CompletionStatus as GeminiCompletionStatus, Error as GeminiError, Gemini, GenerationOptions,
GroundedSearchRequest, MediaInput as GeminiMediaInput, MultimodalRequest, ServiceTier,
StructuredOutput, TextModel, ThinkingLevel, TokenUsage as GeminiTokenUsage,
};
use kcode_openai_api::{
AudioInput, Error as OpenAiError, ImageAnalysisRequest,
ImageAnalysisStatus as OpenAiImageStatus, ImageAnalysisUsage, ImageInput as OpenAiImageInput,
ImageMediaType as OpenAiImageMediaType, OpenAi,
TranscriptionRequest as OpenAiTranscriptionRequest, TranscriptionUsage,
};
use kcode_web_fetch::{ErrorKind as WebFetchErrorKind, WebFetcher};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use tokio::sync::watch;
use uuid::Uuid;
use defaults::*;
pub use error::{Error, ErrorKind, Result};
use receipts::ReceiptStore;
pub use receipts::{DailyUsage, DailyUsageKey, Metering, TokenUsage, UsageReceipt};
pub struct Config {
pub openai_api_key: Option<String>,
pub gemini_api_key: Option<String>,
pub codex_catalog_cache: CatalogCache,
pub receipt_directory: PathBuf,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RuntimeModel {
pub model: String,
pub reasoning_effort: String,
pub context_window_tokens: u64,
pub max_input_tokens: u64,
}
#[derive(Clone)]
pub struct Intelligence {
codex: Codex,
agent: kcode_codex_runtime_v2::Codex,
openai: Option<OpenAi>,
gemini: Option<Gemini>,
web_fetcher: WebFetcher,
document_extractor: DocumentExtractor,
active_operations: ActiveOperations,
receipts: ReceiptStore,
}
#[derive(Clone)]
pub struct UserIntelligence {
service: Intelligence,
user_id: String,
}
#[derive(Clone, Default)]
struct ActiveOperations {
senders: Arc<Mutex<HashMap<Uuid, watch::Sender<bool>>>>,
}
struct ActiveOperation {
id: Uuid,
operations: ActiveOperations,
cancellation: watch::Receiver<bool>,
owns_registration: bool,
}
pub struct AgentTurn {
inner: kcode_codex_runtime_v2::AgentTurn,
operation: ActiveOperation,
user_id: String,
model: String,
receipts: ReceiptStore,
receipt_recorded: bool,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SearchRequest {
pub question: String,
pub model: String,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct WebSource {
pub title: String,
pub url: String,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SearchResponse {
pub answer: String,
pub sources: Vec<WebSource>,
pub model: String,
pub usage: Option<TokenUsage>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct FetchRequest {
pub url: String,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct FetchResponse {
pub url: String,
pub title: Option<String>,
pub content_type: String,
pub content: String,
pub truncated: bool,
pub retrieved_at: DateTime<Utc>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum MediaKind {
Image,
Audio,
Video,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Media {
pub kind: MediaKind,
pub bytes: Vec<u8>,
pub file_name: String,
pub content_type: String,
}
impl Media {
pub fn new(
kind: MediaKind,
bytes: Vec<u8>,
file_name: impl Into<String>,
content_type: impl Into<String>,
) -> Result<Self> {
if bytes.is_empty() || bytes.len() > MAX_MEDIA_ANNOTATION_BYTES {
return Err(Error::invalid(format!(
"media must contain between 1 and {MAX_MEDIA_ANNOTATION_BYTES} bytes"
)));
}
let file_name = file_name.into();
let mut content_type = normalized_content_type(&content_type.into());
if kind == MediaKind::Audio && is_ogg(&file_name, &content_type) {
content_type = "audio/ogg".into();
}
Ok(Self {
kind,
bytes,
file_name,
content_type,
})
}
pub fn audio(
bytes: Vec<u8>,
file_name: impl Into<String>,
content_type: impl Into<String>,
) -> Result<Self> {
Self::new(MediaKind::Audio, bytes, file_name, content_type)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TranscriptionRequest {
pub prompt: String,
pub model: String,
pub media: Media,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct TranscriptionResponse {
pub model: String,
pub text: String,
pub metering: Metering,
}
#[derive(Clone, Debug, PartialEq)]
pub struct StructuredAudioRequest {
pub operation: String,
pub prompt: String,
pub model: String,
pub media: Media,
pub schema: Value,
pub max_output_tokens: u32,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct StructuredAudioResponse {
pub model: String,
pub text: String,
pub usage: TokenUsage,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TextGenerationRequest {
pub operation: String,
pub prompt: String,
pub model: String,
pub reasoning_effort: ReasoningEffort,
pub timeout: Duration,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct TextGenerationResponse {
pub model: String,
pub text: String,
pub thread_id: String,
pub usage: Option<TokenUsage>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct AnnotationRequest {
pub prompt: String,
pub model: String,
pub media: Media,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AnnotationResponse {
pub complete: bool,
pub model: String,
pub file_name: String,
pub content_type: String,
pub text: String,
pub incomplete_reason: Option<String>,
pub usage: Option<TokenUsage>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Document {
pub bytes: Vec<u8>,
pub file_name: String,
pub content_type: String,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct DocumentExtraction {
pub file_name: String,
pub content_type: String,
pub format: String,
pub text: String,
pub characters: usize,
pub truncated: bool,
}
pub async fn open(config: Config) -> anyhow::Result<(Intelligence, RuntimeModel)> {
let mut codex_config = CodexConfig::new(DEFAULT_MODEL);
codex_config.base_instruction = KENNEDY_CODEX_BASE_INSTRUCTION.into();
codex_config.validation_reasoning_effort = GENERATION_REASONING_EFFORT;
let codex = Codex::open(codex_config, config.codex_catalog_cache)
.await
.context("opening Kennedy Codex runtime")?;
let limits = codex
.catalog()
.model_limits(DEFAULT_MODEL)
.with_context(|| format!("Codex model {DEFAULT_MODEL} is absent from the catalog"))?;
let runtime = RuntimeModel {
model: DEFAULT_MODEL.into(),
reasoning_effort: GENERATION_REASONING_EFFORT.as_str().into(),
context_window_tokens: limits.context_window_tokens(),
max_input_tokens: limits.max_input_tokens(),
};
let agent_config = kcode_codex_runtime_v2::CodexConfig {
executable: codex.catalog().executable().to_owned(),
base_instruction: KENNEDY_CODEX_BASE_INSTRUCTION.into(),
model_catalog: Some(codex.catalog().path().to_owned()),
..kcode_codex_runtime_v2::CodexConfig::default()
};
let agent = kcode_codex_runtime_v2::Codex::open(agent_config)
.await
.context("opening Kennedy Codex runtime v2")?;
let openai = config
.openai_api_key
.filter(|value| !value.trim().is_empty())
.map(OpenAi::open)
.transpose()
.context("opening OpenAI client")?;
let gemini = config
.gemini_api_key
.filter(|value| !value.trim().is_empty())
.map(Gemini::open)
.transpose()
.context("opening Gemini client")?;
Ok((
Intelligence {
codex,
agent,
openai,
gemini,
web_fetcher: WebFetcher::default(),
document_extractor: DocumentExtractor::default(),
active_operations: ActiveOperations::default(),
receipts: ReceiptStore::open(config.receipt_directory)
.map_err(anyhow::Error::new)
.context("opening intelligence usage receipts")?,
},
runtime,
))
}
impl Intelligence {
pub fn for_user(&self, user_id: impl Into<String>) -> Result<UserIntelligence> {
let user_id = user_id.into();
if user_id.trim().is_empty() || user_id.chars().count() > 256 {
return Err(Error::invalid(
"user_id must contain between 1 and 256 characters",
));
}
Ok(UserIntelligence {
service: self.clone(),
user_id,
})
}
pub fn cancel(&self, operation_id: Uuid) -> Result<bool> {
self.active_operations.cancel(operation_id)
}
pub fn receipts(&self) -> Result<Vec<UsageReceipt>> {
self.receipts.receipts()
}
pub fn daily_usage(
&self,
day: NaiveDate,
) -> Result<std::collections::BTreeMap<DailyUsageKey, DailyUsage>> {
self.receipts.daily_usage(day)
}
pub async fn extract_document(&self, document: Document) -> Result<DocumentExtraction> {
extract_document(&self.document_extractor, document).await
}
}
impl UserIntelligence {
pub fn user_id(&self) -> &str {
&self.user_id
}
pub async fn start_agent_turn(
&self,
operation_id: Uuid,
request: kcode_codex_runtime_v2::AgentRequest,
) -> Result<AgentTurn> {
validate_model(&request.model)?;
if request.previous_thread_id.is_some() {
return Err(Error::invalid(
"agent turns must submit complete input without a previous thread; this keeps each usage receipt scoped to one call",
));
}
let operation = self.service.active_operations.register(operation_id)?;
let model = request.model.clone();
let inner = self.account_result(
"agent_turn",
&model,
self.service
.agent
.start_turn(request)
.await
.map_err(codex_v2_error),
)?;
Ok(AgentTurn {
inner,
operation,
user_id: self.user_id.clone(),
model,
receipts: self.service.receipts.clone(),
receipt_recorded: false,
})
}
pub async fn search(&self, request: SearchRequest) -> Result<SearchResponse> {
let question = request.question.trim();
if question.is_empty() || question.chars().count() > 4_000 {
return Err(Error::invalid(
"question must contain between 1 and 4000 characters",
));
}
validate_model(&request.model)?;
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let started = Instant::now();
let response = if let Some(model) = gemini_model(&request.model) {
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"Gemini search is not configured.",
)
})?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
FAST_SEARCH_TIMEOUT,
gemini.grounded_search_with_model(
model,
GroundedSearchRequest::new(question),
),
) => result
.map_err(|_| Error::provider("provider_timeout", "Gemini search timed out."))
.and_then(|result| result.map_err(gemini_error)),
};
let result = self.account_result("web_search", &request.model, result)?;
let interaction = result.interaction;
let usage = gemini_usage(&interaction.usage);
self.record_tokens(
"web_search",
&request.model,
&interaction.model,
Some(usage),
Some(interaction.id.clone()),
None,
)?;
SearchResponse {
answer: interaction
.text
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| {
Error::provider("provider_error", "Gemini search returned no answer text.")
})?,
sources: result
.sources
.into_iter()
.map(|source| WebSource {
title: source.title,
url: source.url,
})
.collect(),
model: interaction.model,
usage: Some(usage),
}
} else {
let (reasoning, context, depth, timeout) = codex_search_profile(&request.model)?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = self.service.codex.web_search(CodexSearchRequest {
question: question.to_owned(),
model: request.model.clone(),
reasoning_effort: reasoning,
context,
depth,
timeout,
}) => result.map_err(codex_error),
};
let result = self.account_result("web_search", &request.model, result)?;
let usage = result.usage.as_ref().map(codex_usage);
self.record_tokens(
"web_search",
&request.model,
&request.model,
usage,
None,
None,
)?;
SearchResponse {
answer: result.answer,
sources: result
.sources
.into_iter()
.map(|source| WebSource {
title: source.title,
url: source.url,
})
.collect(),
model: request.model.clone(),
usage,
}
};
tracing::info!(
user_id = %self.user_id,
operation_id = %request.operation_id,
model = %response.model,
duration_ms = started.elapsed().as_millis(),
"Intelligence search completed"
);
Ok(response)
}
pub async fn transcribe_structured_audio(
&self,
request: StructuredAudioRequest,
) -> Result<StructuredAudioResponse> {
validate_operation(&request.operation)?;
validate_prompt(&request.prompt)?;
validate_model(&request.model)?;
if request.media.kind != MediaKind::Audio {
return Err(Error::invalid(
"structured audio inference requires audio media",
));
}
if request.max_output_tokens == 0 || request.max_output_tokens > 65_536 {
return Err(Error::invalid(
"max_output_tokens must be between 1 and 65536",
));
}
let model = gemini_model(&request.model).ok_or_else(|| {
Error::invalid("structured audio inference requires a supported exact Gemini model")
})?;
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"Gemini structured audio inference is not configured.",
)
})?;
let media = gemini_media(&request.media)?;
let structured_output = StructuredOutput::new(request.schema).map_err(gemini_error)?;
let mut provider_request = MultimodalRequest::new(request.prompt, vec![media]);
provider_request.options = GenerationOptions {
max_output_tokens: Some(request.max_output_tokens),
temperature: None,
thinking_level: Some(ThinkingLevel::High),
service_tier: ServiceTier::Standard,
};
provider_request.structured_output = Some(structured_output);
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
MEDIA_ANNOTATION_TIMEOUT,
gemini.infer_multimodal(model, provider_request),
) => result
.map_err(|_| Error::provider("provider_timeout", "Gemini structured audio inference timed out."))
.and_then(|result| result.map_err(gemini_error)),
};
let result = self.account_result(&request.operation, &request.model, result)?;
let usage = gemini_usage(&result.usage);
self.record_tokens(
&request.operation,
&request.model,
&result.model,
Some(usage),
Some(result.id.clone()),
None,
)?;
if result.status != GeminiCompletionStatus::Completed {
return Err(Error::provider(
"provider_incomplete",
"Gemini structured audio inference did not complete.",
));
}
let text = result
.text
.filter(|text| !text.trim().is_empty())
.ok_or_else(|| {
Error::provider(
"provider_empty_output",
"Gemini structured audio inference returned no text.",
)
})?;
Ok(StructuredAudioResponse {
model: result.model,
text,
usage,
})
}
pub async fn generate_text(
&self,
request: TextGenerationRequest,
) -> Result<TextGenerationResponse> {
validate_operation(&request.operation)?;
validate_model(&request.model)?;
if request.prompt.trim().is_empty() || request.prompt.chars().count() > 1_000_000 {
return Err(Error::invalid(
"generation prompt must contain between 1 and 1000000 characters",
));
}
if request.timeout.is_zero() || request.timeout > Duration::from_secs(60 * 60) {
return Err(Error::invalid(
"generation timeout must be between 1 second and 1 hour",
));
}
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let mut provider_request =
CodexGenerationRequest::new(request.prompt, request.model.clone());
provider_request.reasoning_effort = request.reasoning_effort;
provider_request.ephemeral = true;
provider_request.timeout = request.timeout;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = self.service.codex.generate(provider_request) => result.map_err(codex_error),
};
let result = self.account_result(&request.operation, &request.model, result)?;
let usage = result.usage.as_ref().map(codex_usage);
self.record_tokens(
&request.operation,
&request.model,
&request.model,
usage,
None,
Some(result.thread_id.clone()),
)?;
Ok(TextGenerationResponse {
model: request.model,
text: result.answer,
thread_id: result.thread_id,
usage,
})
}
pub async fn fetch(&self, request: FetchRequest) -> Result<FetchResponse> {
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let fetched = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = self.service.web_fetcher.fetch(&request.url) => result.map_err(web_fetch_error),
}?;
Ok(FetchResponse {
url: fetched.url,
title: fetched.title,
content_type: fetched.content_type,
content: fetched.content,
truncated: fetched.truncated,
retrieved_at: DateTime::<Utc>::from(fetched.retrieved_at),
})
}
pub async fn transcribe(&self, request: TranscriptionRequest) -> Result<TranscriptionResponse> {
validate_prompt(&request.prompt)?;
if request.media.kind != MediaKind::Audio {
return Err(Error::invalid("transcription requires audio media"));
}
validate_model(&request.model)?;
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let response = if request.model == kcode_openai_api::GPT_4O_TRANSCRIBE {
let openai = self.service.openai.as_ref().ok_or_else(|| {
Error::unavailable(
"transcription_unavailable",
"OpenAI audio transcription is not configured.",
)
})?;
let input = AudioInput::new(
safe_audio_filename(&request.media.file_name, &request.media.content_type),
request.media.content_type.clone(),
request.media.bytes,
)
.map_err(openai_error)?;
let mut provider_request = OpenAiTranscriptionRequest::new(input);
provider_request.prompt = Some(request.prompt);
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
MEDIA_ANNOTATION_TIMEOUT,
openai.transcribe(provider_request),
) => result
.map_err(|_| Error::provider("provider_timeout", "OpenAI transcription timed out."))
.and_then(|result| result.map_err(openai_error)),
};
let result = self.account_result("transcribe_audio", &request.model, result)?;
let metering = result
.usage
.map(transcription_usage)
.unwrap_or(Metering::Unavailable);
self.record_metering(
"transcribe_audio",
&request.model,
&request.model,
metering.clone(),
None,
None,
)?;
TranscriptionResponse {
model: request.model.clone(),
text: result.text,
metering,
}
} else {
let model = gemini_model(&request.model).ok_or_else(|| {
Error::invalid("transcription model must be gpt-4o-transcribe or a supported exact Gemini model")
})?;
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"Gemini transcription is not configured.",
)
})?;
let media = gemini_media(&request.media)?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
MEDIA_ANNOTATION_TIMEOUT,
gemini.infer_multimodal(
model,
MultimodalRequest::new(request.prompt, vec![media]),
),
) => result
.map_err(|_| Error::provider("provider_timeout", "Gemini transcription timed out."))
.and_then(|result| result.map_err(gemini_error)),
};
let result = self.account_result("transcribe_audio", &request.model, result)?;
let metering = Metering::Tokens(gemini_usage(&result.usage));
self.record_metering(
"transcribe_audio",
&request.model,
&result.model,
metering.clone(),
Some(result.id.clone()),
None,
)?;
TranscriptionResponse {
model: result.model,
text: result
.text
.filter(|text| !text.trim().is_empty())
.ok_or_else(|| {
Error::provider(
"empty_transcription",
"Gemini returned no transcription text.",
)
})?,
metering,
}
};
Ok(response)
}
pub async fn annotate(&self, request: AnnotationRequest) -> Result<AnnotationResponse> {
validate_prompt(&request.prompt)?;
validate_model(&request.model)?;
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let file_name = request.media.file_name.clone();
let content_type = request.media.content_type.clone();
let response = if let Some(model) = gemini_model(&request.model) {
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"Gemini media annotation is not configured.",
)
})?;
let media = gemini_media(&request.media)?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
MEDIA_ANNOTATION_TIMEOUT,
gemini.infer_multimodal(
model,
MultimodalRequest::new(request.prompt, vec![media]),
),
) => result
.map_err(|_| Error::provider("provider_timeout", "Gemini annotation timed out."))
.and_then(|result| result.map_err(gemini_error)),
};
let result = self.account_result("annotate_media", &request.model, result)?;
let usage = gemini_usage(&result.usage);
self.record_tokens(
"annotate_media",
&request.model,
&result.model,
Some(usage),
Some(result.id.clone()),
None,
)?;
AnnotationResponse {
complete: result.status == GeminiCompletionStatus::Completed,
model: result.model,
file_name,
content_type,
text: result
.text
.filter(|text| !text.trim().is_empty())
.ok_or_else(|| {
Error::provider("empty_annotation", "Gemini returned no annotation text.")
})?,
incomplete_reason: None,
usage: Some(usage),
}
} else if request.model == "gpt-5.6" {
if request.media.kind != MediaKind::Image {
return Err(Error::invalid(
"OpenAI media annotation accepts images only",
));
}
let openai = self.service.openai.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"OpenAI media annotation is not configured.",
)
})?;
let image = OpenAiImageInput::new(
openai_image_media_type(&request.media.content_type)?,
request.media.bytes,
)
.map_err(openai_error)?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
MEDIA_ANNOTATION_TIMEOUT,
openai.analyze_image(ImageAnalysisRequest::new(image, request.prompt)),
) => result
.map_err(|_| Error::provider("provider_timeout", "OpenAI annotation timed out."))
.and_then(|result| result.map_err(openai_error)),
};
let result = self.account_result("annotate_media", &request.model, result)?;
let usage = result.usage.as_ref().map(openai_image_usage);
self.record_tokens(
"annotate_media",
&request.model,
&result.model,
usage,
None,
None,
)?;
let (complete, incomplete_reason) = match result.status {
OpenAiImageStatus::Completed => (true, None),
OpenAiImageStatus::Incomplete { reason } => (false, reason),
};
AnnotationResponse {
complete,
model: result.model,
file_name,
content_type,
text: result.text,
incomplete_reason,
usage,
}
} else if matches!(
request.model.as_str(),
"gpt-5.6-sol" | "gpt-5.6-terra" | "gpt-5.6-luna"
) {
if request.media.kind != MediaKind::Image {
return Err(Error::invalid("Codex media annotation accepts images only"));
}
let image = kcode_codex_runtime_v2::ImageInput::new(
codex_image_media_type(&request.media.content_type)?,
request.media.bytes,
)
.map_err(codex_v2_error)?;
let turn_result = self
.service
.agent
.start_image_turn(kcode_codex_runtime_v2::ImageTurnRequest::new(
request.prompt,
request.model.clone(),
vec![image],
))
.await
.map_err(codex_v2_error);
let mut turn = self.account_result("annotate_media", &request.model, turn_result)?;
let completed = loop {
let event = tokio::select! {
_ = operation.cancelled() => {
turn.cancel();
self.record_metering(
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(Error::cancelled());
}
event = turn.next_event() => event,
};
match event {
Some(Ok(kcode_codex_runtime_v2::AgentEvent::ProviderInput(_))) => {}
Some(Ok(kcode_codex_runtime_v2::AgentEvent::ToolCall(_))) => {
turn.cancel();
self.record_metering(
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(Error::provider(
"provider_error",
"A tool-free Codex image turn requested a tool.",
));
}
Some(Ok(kcode_codex_runtime_v2::AgentEvent::Completed(completed))) => {
break completed;
}
Some(Err(error)) => {
self.record_metering(
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(codex_v2_error(error));
}
None => {
self.record_metering(
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(Error::provider(
"empty_annotation",
"Codex ended without annotation text.",
));
}
}
};
let usage = completed.usage.as_ref().map(codex_v2_usage);
self.record_tokens(
"annotate_media",
&request.model,
&request.model,
usage,
Some(completed.turn_id.clone()),
Some(completed.thread_id.clone()),
)?;
AnnotationResponse {
complete: true,
model: request.model.clone(),
file_name,
content_type,
text: completed.answer,
incomplete_reason: None,
usage,
}
} else {
return Err(Error::invalid(format!(
"unsupported exact annotation model {}",
request.model
)));
};
Ok(response)
}
fn record_tokens(
&self,
operation: &str,
requested_model: &str,
actual_model: &str,
usage: Option<TokenUsage>,
provider_request_id: Option<String>,
provider_thread_id: Option<String>,
) -> Result<()> {
self.record_metering(
operation,
requested_model,
actual_model,
usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
provider_request_id,
provider_thread_id,
)
}
fn record_metering(
&self,
operation: &str,
requested_model: &str,
actual_model: &str,
metering: Metering,
provider_request_id: Option<String>,
provider_thread_id: Option<String>,
) -> Result<()> {
let mut receipt = UsageReceipt::new(
self.user_id.clone(),
operation,
requested_model,
actual_model,
metering,
);
receipt.provider_request_id = provider_request_id;
receipt.provider_thread_id = provider_thread_id;
self.service.receipts.record(&receipt)
}
fn account_result<T>(
&self,
operation: &str,
requested_model: &str,
result: Result<T>,
) -> Result<T> {
match result {
Ok(value) => Ok(value),
Err(error) => {
self.record_metering(
operation,
requested_model,
requested_model,
Metering::Unavailable,
None,
None,
)?;
Err(error)
}
}
}
}
async fn extract_document(
extractor: &DocumentExtractor,
document: Document,
) -> Result<DocumentExtraction> {
let extractor = extractor.clone();
let extracted = tokio::task::spawn_blocking(move || {
extractor.extract(DocumentInput {
file_name: document.file_name,
content_type: document.content_type,
data: document.bytes,
})
})
.await
.map_err(|_| {
Error::internal(
"document_extraction_failed",
"The document extraction worker stopped unexpectedly.",
)
})?
.map_err(document_error)?;
Ok(DocumentExtraction {
file_name: extracted.file_name,
content_type: extracted.content_type,
format: extracted.format.as_str().into(),
text: extracted.text,
characters: extracted.characters,
truncated: extracted.truncated,
})
}
impl AgentTurn {
pub async fn next_event(&mut self) -> Result<Option<kcode_codex_runtime_v2::AgentEvent>> {
let event = tokio::select! {
_ = self.operation.cancelled() => {
self.inner.cancel();
self.record_unavailable()?;
return Err(Error::cancelled());
}
event = self.inner.next_event() => match event {
Some(Ok(event)) => Some(event),
Some(Err(error)) => {
self.record_unavailable()?;
return Err(codex_v2_error(error));
}
None => {
self.record_unavailable()?;
None
},
}
};
if let Some(kcode_codex_runtime_v2::AgentEvent::Completed(completed)) = &event
&& !self.receipt_recorded
{
let mut receipt = UsageReceipt::new(
self.user_id.clone(),
"agent_turn",
self.model.clone(),
self.model.clone(),
completed
.usage
.as_ref()
.map(codex_v2_usage)
.map(Metering::Tokens)
.unwrap_or(Metering::Unavailable),
);
receipt.provider_request_id = Some(completed.turn_id.clone());
receipt.provider_thread_id = Some(completed.thread_id.clone());
self.receipts.record(&receipt)?;
self.receipt_recorded = true;
}
Ok(event)
}
fn record_unavailable(&mut self) -> Result<()> {
if self.receipt_recorded {
return Ok(());
}
let receipt = UsageReceipt::new(
self.user_id.clone(),
"agent_turn",
self.model.clone(),
self.model.clone(),
Metering::Unavailable,
);
self.receipts.record(&receipt)?;
self.receipt_recorded = true;
Ok(())
}
pub async fn respond(
&self,
call_id: &str,
result: kcode_codex_runtime_v2::ToolResult,
) -> Result<()> {
self.inner
.respond(call_id, result)
.await
.map_err(codex_v2_error)
}
}
impl Drop for AgentTurn {
fn drop(&mut self) {
let _ = self.record_unavailable();
}
}
impl ActiveOperations {
fn register(&self, id: Uuid) -> Result<ActiveOperation> {
let (sender, cancellation) = watch::channel(false);
let mut senders = self.senders.lock().map_err(|_| {
Error::internal(
"operation_registry_unavailable",
"The operation registry is unavailable.",
)
})?;
if senders.contains_key(&id) {
return Err(Error::conflict(
"operation_in_progress",
"An operation with this identifier is already running.",
));
}
senders.insert(id, sender);
Ok(ActiveOperation {
id,
operations: self.clone(),
cancellation,
owns_registration: true,
})
}
fn register_request(
&self,
request_id: Uuid,
parent_operation_id: Option<Uuid>,
) -> Result<ActiveOperation> {
if let Some(parent_id) = parent_operation_id {
if request_id == parent_id {
return Err(Error::invalid(
"operation_id and parent_operation_id must be different",
));
}
let cancellation = self
.senders
.lock()
.map_err(|_| {
Error::internal(
"operation_registry_unavailable",
"The operation registry is unavailable.",
)
})?
.get(&parent_id)
.map(watch::Sender::subscribe)
.ok_or_else(|| {
Error::conflict(
"parent_operation_not_running",
"The parent operation is no longer running.",
)
})?;
Ok(ActiveOperation {
id: parent_id,
operations: self.clone(),
cancellation,
owns_registration: false,
})
} else {
self.register(request_id)
}
}
fn cancel(&self, id: Uuid) -> Result<bool> {
let sender = self
.senders
.lock()
.map_err(|_| {
Error::internal(
"operation_registry_unavailable",
"The operation registry is unavailable.",
)
})?
.get(&id)
.cloned();
Ok(sender.is_some_and(|sender| sender.send(true).is_ok()))
}
fn remove(&self, id: Uuid) {
if let Ok(mut senders) = self.senders.lock() {
senders.remove(&id);
}
}
}
impl ActiveOperation {
async fn cancelled(&mut self) {
if *self.cancellation.borrow() {
return;
}
while self.cancellation.changed().await.is_ok() {
if *self.cancellation.borrow() {
return;
}
}
}
}
impl Drop for ActiveOperation {
fn drop(&mut self) {
if self.owns_registration {
self.operations.remove(self.id);
}
}
}
fn validate_model(model: &str) -> Result<()> {
if model.trim().is_empty()
|| model.chars().count() > 128
|| !model
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_'))
{
return Err(Error::invalid(
"model must be an exact safe model identifier",
));
}
Ok(())
}
fn validate_operation(operation: &str) -> Result<()> {
if operation.is_empty()
|| operation.len() > 64
|| !operation
.bytes()
.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
{
return Err(Error::invalid(
"operation must be a lowercase identifier of at most 64 bytes",
));
}
Ok(())
}
fn validate_prompt(prompt: &str) -> Result<()> {
if prompt.trim().is_empty() || prompt.chars().count() > MAX_MEDIA_ANNOTATION_PROMPT_CHARACTERS {
return Err(Error::invalid(format!(
"prompt must contain between 1 and {MAX_MEDIA_ANNOTATION_PROMPT_CHARACTERS} characters"
)));
}
Ok(())
}
fn gemini_model(model: &str) -> Option<TextModel> {
match model {
kcode_gemini_api::GEMINI_25_FLASH => Some(TextModel::Flash25),
kcode_gemini_api::GEMINI_31_FLASH_LITE => Some(TextModel::FlashLite),
kcode_gemini_api::GEMINI_31_PRO => Some(TextModel::Pro),
_ => None,
}
}
fn codex_search_profile(
model: &str,
) -> Result<(
kcode_codex_runtime::ReasoningEffort,
kcode_codex_runtime::WebSearchContext,
kcode_codex_runtime::SearchDepth,
Duration,
)> {
match model {
QUALITY_SEARCH_MODEL => Ok((
QUALITY_SEARCH_REASONING,
QUALITY_SEARCH_CONTEXT,
QUALITY_SEARCH_DEPTH,
QUALITY_SEARCH_TIMEOUT,
)),
BALANCED_SEARCH_MODEL => Ok((
BALANCED_SEARCH_REASONING,
BALANCED_SEARCH_CONTEXT,
BALANCED_SEARCH_DEPTH,
BALANCED_SEARCH_TIMEOUT,
)),
_ => Err(Error::invalid(
"unsupported exact web-search model; use a supported Gemini model, gpt-5.6-sol, or gpt-5.6-terra",
)),
}
}
fn normalized_content_type(value: &str) -> String {
value
.split(';')
.next()
.unwrap_or("application/octet-stream")
.trim()
.to_ascii_lowercase()
}
fn is_ogg(file_name: &str, content_type: &str) -> bool {
matches!(content_type, "audio/ogg" | "video/ogg" | "application/ogg")
|| file_name.rsplit_once('.').is_some_and(|(_, extension)| {
matches!(
extension.to_ascii_lowercase().as_str(),
"ogg" | "oga" | "opus"
)
})
}
fn safe_audio_filename(value: &str, content_type: &str) -> String {
let extension = match content_type {
"audio/ogg" | "audio/opus" | "application/ogg" | "video/ogg" => "ogg",
"audio/wav" | "audio/x-wav" => "wav",
"audio/mpeg" | "audio/mp3" => "mp3",
"audio/mp4" => "mp4",
"audio/webm" => "webm",
"audio/flac" | "audio/x-flac" => "flac",
"audio/m4a" => "m4a",
_ => "audio",
};
let cleaned = value
.chars()
.filter(|character| {
character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_')
})
.take(120)
.collect::<String>();
let supported = cleaned.rsplit_once('.').is_some_and(|(_, extension)| {
matches!(
extension.to_ascii_lowercase().as_str(),
"flac"
| "mp3"
| "mp4"
| "mpeg"
| "mpga"
| "m4a"
| "ogg"
| "oga"
| "opus"
| "wav"
| "webm"
)
});
if cleaned.is_empty() || !supported {
format!("voice-note.{extension}")
} else {
cleaned
}
}
fn gemini_media(media: &Media) -> Result<GeminiMediaInput> {
match media.kind {
MediaKind::Image => GeminiMediaInput::image(&media.content_type, media.bytes.clone()),
MediaKind::Audio => GeminiMediaInput::audio(&media.content_type, media.bytes.clone()),
MediaKind::Video => GeminiMediaInput::video(&media.content_type, media.bytes.clone()),
}
.map_err(gemini_error)
}
fn openai_image_media_type(content_type: &str) -> Result<OpenAiImageMediaType> {
match content_type {
"image/png" => Ok(OpenAiImageMediaType::Png),
"image/jpeg" | "image/jpg" => Ok(OpenAiImageMediaType::Jpeg),
"image/webp" => Ok(OpenAiImageMediaType::WebP),
"image/gif" => Ok(OpenAiImageMediaType::Gif),
_ => Err(Error::invalid(
"OpenAI annotations require PNG, JPEG, WebP, or GIF",
)),
}
}
fn codex_image_media_type(content_type: &str) -> Result<kcode_codex_runtime_v2::ImageMediaType> {
match content_type {
"image/png" => Ok(kcode_codex_runtime_v2::ImageMediaType::Png),
"image/jpeg" | "image/jpg" => Ok(kcode_codex_runtime_v2::ImageMediaType::Jpeg),
"image/webp" => Ok(kcode_codex_runtime_v2::ImageMediaType::Webp),
_ => Err(Error::invalid(
"Codex annotations require PNG, JPEG, or WebP",
)),
}
}
fn gemini_usage(usage: &GeminiTokenUsage) -> TokenUsage {
TokenUsage {
input_tokens: usage.input_tokens.saturating_sub(usage.cached_tokens),
cached_input_tokens: usage.cached_tokens,
thinking_tokens: usage.thought_tokens,
output_tokens: usage.output_tokens,
}
}
fn codex_usage(usage: &CodexTokenUsage) -> TokenUsage {
TokenUsage {
input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
cached_input_tokens: usage.cached_input_tokens,
thinking_tokens: usage.reasoning_output_tokens,
output_tokens: usage
.output_tokens
.saturating_sub(usage.reasoning_output_tokens),
}
}
fn codex_v2_usage(usage: &kcode_codex_runtime_v2::TokenUsage) -> TokenUsage {
TokenUsage {
input_tokens: usage.input_tokens.saturating_sub(usage.cached_input_tokens),
cached_input_tokens: usage.cached_input_tokens,
thinking_tokens: usage.reasoning_output_tokens,
output_tokens: usage
.output_tokens
.saturating_sub(usage.reasoning_output_tokens),
}
}
fn openai_image_usage(usage: &ImageAnalysisUsage) -> TokenUsage {
let cached = usage.cached_input_tokens.unwrap_or(0);
let thinking = usage.reasoning_output_tokens.unwrap_or(0);
TokenUsage {
input_tokens: usage.input_tokens.saturating_sub(cached),
cached_input_tokens: cached,
thinking_tokens: thinking,
output_tokens: usage.output_tokens.saturating_sub(thinking),
}
}
fn transcription_usage(usage: TranscriptionUsage) -> Metering {
match usage {
TranscriptionUsage::DurationSeconds(seconds) => Metering::DurationSeconds { seconds },
TranscriptionUsage::Tokens(tokens) => Metering::Tokens(TokenUsage {
input_tokens: tokens.input_tokens,
cached_input_tokens: 0,
thinking_tokens: 0,
output_tokens: tokens.output_tokens,
}),
}
}
fn codex_error(error: kcode_codex_runtime::Error) -> Error {
match error.kind() {
CodexErrorKind::InvalidInput => Error::invalid(error.message()),
CodexErrorKind::Authentication => {
Error::unavailable("provider_not_configured", error.message())
}
CodexErrorKind::Unavailable => Error::unavailable("provider_unavailable", error.message()),
CodexErrorKind::RateLimited | CodexErrorKind::Capacity => {
Error::unavailable("provider_rate_limited", error.message())
}
CodexErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
CodexErrorKind::InputTooLarge => Error::invalid(error.message()),
CodexErrorKind::EmptyOutput | CodexErrorKind::Protocol => {
Error::provider("provider_error", error.message())
}
}
}
fn codex_v2_error(error: kcode_codex_runtime_v2::Error) -> Error {
use kcode_codex_runtime_v2::ErrorKind;
match error.kind() {
ErrorKind::InvalidInput => Error::invalid(error.message()),
ErrorKind::Unavailable | ErrorKind::Authentication => {
Error::unavailable("provider_unavailable", error.message())
}
ErrorKind::Timeout => Error::provider("provider_timeout", error.message()),
ErrorKind::Protocol => Error::provider("provider_error", error.message()),
ErrorKind::Cancelled => Error::cancelled(),
}
}
fn gemini_error(error: GeminiError) -> Error {
match &error {
GeminiError::InvalidApiKey => {
Error::unavailable("provider_not_configured", error.to_string())
}
GeminiError::InvalidInput(_) => Error::invalid(error.to_string()),
GeminiError::SpendingLimitReached { .. } => {
Error::unavailable("provider_rate_limited", error.to_string())
}
GeminiError::Accounting(_)
| GeminiError::Transport(_)
| GeminiError::Provider { .. }
| GeminiError::Protocol(_) => Error::provider("provider_error", error.to_string()),
}
}
fn openai_error(error: OpenAiError) -> Error {
match &error {
OpenAiError::InvalidApiKey => {
Error::unavailable("provider_not_configured", error.to_string())
}
OpenAiError::InvalidInput(_) => Error::invalid(error.to_string()),
OpenAiError::Transport(_) | OpenAiError::Provider { .. } | OpenAiError::Protocol(_) => {
Error::provider("provider_error", error.to_string())
}
}
}
fn web_fetch_error(error: kcode_web_fetch::Error) -> Error {
match error.kind() {
WebFetchErrorKind::InvalidInput | WebFetchErrorKind::UnsafeDestination => {
Error::invalid(error.message())
}
WebFetchErrorKind::Timeout => Error::provider("web_fetch_timeout", error.message()),
WebFetchErrorKind::UnsupportedContent => {
Error::invalid(format!("unsupported web content: {}", error.message()))
}
WebFetchErrorKind::Transport
| WebFetchErrorKind::HttpStatus
| WebFetchErrorKind::EmptyContent => Error::provider("web_fetch_failed", error.message()),
}
}
fn document_error(error: kcode_doc_extraction::Error) -> Error {
match error.kind() {
DocumentErrorKind::InvalidInput | DocumentErrorKind::UnsupportedFormat => {
Error::invalid(error.message())
}
DocumentErrorKind::ExtractionFailed | DocumentErrorKind::EmptyText => {
Error::provider("document_extraction_failed", error.message())
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn audio_kind_overrides_mislabeled_ogg_video_mime() {
let media = Media::audio(vec![1], "voice.ogg", "video/ogg").unwrap();
assert_eq!(media.kind, MediaKind::Audio);
assert_eq!(media.content_type, "audio/ogg");
}
#[test]
fn exact_search_models_replace_modes() {
assert!(codex_search_profile("gpt-5.6-sol").is_ok());
assert!(codex_search_profile("gpt-5.6-terra").is_ok());
assert!(codex_search_profile("fast").is_err());
}
}