#![forbid(unsafe_code)]
mod defaults;
mod error;
mod pricing;
mod receipts;
use std::{
collections::{HashMap, VecDeque},
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::{
AgentTool as GeminiAgentTool, AgentTurnRequest as GeminiAgentRequest,
CompletionStatus as GeminiCompletionStatus, Error as GeminiError, Gemini, GenerationOptions,
GroundedSearchRequest, MediaInput as GeminiMediaInput, MultimodalRequest, NanoBananaProRequest,
ServiceTier, StructuredOutput, TextModel, ThinkingLevel, TokenUsage as GeminiTokenUsage,
};
use kcode_openai_api::{
AgentTool as OpenAiAgentTool, AgentTurnRequest as OpenAiAgentRequest, AudioInput,
Error as OpenAiError, ImageAnalysisRequest, ImageAnalysisStatus as OpenAiImageStatus,
ImageAnalysisUsage, ImageEditRequest as OpenAiImageEditRequest,
ImageGenerationRequest as OpenAiImageRequest, ImageInput as OpenAiImageInput,
ImageMediaType as OpenAiImageMediaType, ImageUsage as OpenAiGenerationUsage, OpenAi,
TranscriptionRequest as OpenAiTranscriptionRequest, TranscriptionUsage,
};
use kcode_web_fetch::{ErrorKind as WebFetchErrorKind, WebFetcher};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use tokio::sync::watch;
use uuid::Uuid;
use defaults::*;
pub use error::{Error, ErrorKind, Result};
pub use pricing::{
CostAccuracy, CostEstimate, PRICING_VERSION, estimate_cost, estimate_token_cost,
};
use receipts::ReceiptStore;
pub use receipts::{DailyUsage, DailyUsageKey, Metering, TokenUsage, UsageReceipt};
#[derive(Clone, Debug, PartialEq)]
pub struct Accounted<T> {
pub value: T,
pub receipt: UsageReceipt,
}
impl<T> Accounted<T> {
pub fn map<U>(self, map: impl FnOnce(T) -> U) -> Accounted<U> {
Accounted {
value: map(self.value),
receipt: self.receipt,
}
}
}
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, Copy, Debug, Eq, PartialEq)]
pub enum AgentProvider {
Codex,
OpenAi,
Gemini,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ResolvedAgentModel {
pub requested_model: String,
pub provider_model: String,
pub provider: AgentProvider,
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,
model_cache: Arc<Mutex<HashMap<String, ResolvedAgentModel>>>,
}
#[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>,
parent_cancellation: Option<watch::Receiver<bool>>,
}
enum AgentTurnBackend {
Codex(kcode_codex_runtime_v2::AgentTurn),
Buffered(BufferedAgentTurn),
}
struct BufferedAgentTurn {
events: VecDeque<kcode_codex_runtime_v2::AgentEvent>,
pending_call_id: Option<String>,
completed: Option<kcode_codex_runtime_v2::CompletedTurn>,
}
pub struct AgentTurn {
inner: AgentTurnBackend,
operation: ActiveOperation,
user_id: String,
requested_model: String,
actual_model: String,
provider_request_id: Option<String>,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
receipts: ReceiptStore,
receipt: Option<UsageReceipt>,
}
#[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>,
pub cost: Option<CostEstimate>,
}
#[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,
pub cost: Option<CostEstimate>,
}
#[derive(Clone, Debug, PartialEq)]
pub struct AudioAnalysisRequest {
pub operation: String,
pub prompt: String,
pub model: String,
pub media: Media,
pub schema: Option<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 AudioAnalysisResponse {
pub model: String,
pub text: String,
pub usage: TokenUsage,
pub cost: CostEstimate,
}
#[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,
pub cost: CostEstimate,
}
#[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>,
pub cost: Option<CostEstimate>,
}
#[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>,
pub cost: Option<CostEstimate>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ImageRequest {
pub model: String,
pub prompt: String,
pub references: Vec<Media>,
pub operation_id: Uuid,
pub parent_operation_id: Option<Uuid>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ImageResponse {
pub model: String,
pub content_type: String,
pub bytes: Vec<u8>,
pub usage: Option<TokenUsage>,
pub cost: Option<CostEstimate>,
}
#[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")?,
model_cache: Arc::new(Mutex::new(HashMap::new())),
},
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
}
pub async fn resolve_agent_model(&self, requested: &str) -> Result<ResolvedAgentModel> {
validate_agent_model(requested)?;
if let Some(resolved) = self
.model_cache
.lock()
.map_err(|_| {
Error::internal("model_cache_unavailable", "The model cache is unavailable.")
})?
.get(requested)
.cloned()
{
return Ok(resolved);
}
let resolved = if let Some(provider_model) = requested.strip_prefix("codex/") {
let limits = self
.codex
.catalog()
.model_limits(provider_model)
.ok_or_else(|| {
Error::invalid(format!(
"{provider_model:?} is not an available model in the Codex catalog"
))
})?;
ResolvedAgentModel {
requested_model: requested.to_owned(),
provider_model: provider_model.to_owned(),
provider: AgentProvider::Codex,
context_window_tokens: limits.context_window_tokens(),
max_input_tokens: limits.max_input_tokens(),
}
} else if requested.ends_with("-sol")
|| requested.ends_with("-terra")
|| requested.ends_with("-luna")
{
let limits = self
.codex
.catalog()
.model_limits(requested)
.ok_or_else(|| {
Error::invalid(format!(
"{requested:?} is not an available model in the Codex catalog"
))
})?;
ResolvedAgentModel {
requested_model: requested.to_owned(),
provider_model: requested.to_owned(),
provider: AgentProvider::Codex,
context_window_tokens: limits.context_window_tokens(),
max_input_tokens: limits.max_input_tokens(),
}
} else if requested.starts_with("gemini-") {
let gemini = self.gemini.as_ref().ok_or_else(|| {
Error::unavailable("provider_not_configured", "Gemini is not configured.")
})?;
let metadata =
tokio::time::timeout(MODEL_DISCOVERY_TIMEOUT, gemini.model_metadata(requested))
.await
.map_err(|_| {
Error::provider("provider_timeout", "Gemini model discovery timed out.")
})?
.map_err(gemini_error)?;
resolved_api_model(
requested,
metadata.id,
AgentProvider::Gemini,
metadata.context_window_tokens,
metadata.max_input_tokens,
)
} else {
let openai = self.openai.as_ref().ok_or_else(|| {
Error::unavailable("provider_not_configured", "OpenAI is not configured.")
})?;
let metadata =
tokio::time::timeout(MODEL_DISCOVERY_TIMEOUT, openai.model_metadata(requested))
.await
.map_err(|_| {
Error::provider("provider_timeout", "OpenAI model discovery timed out.")
})?
.map_err(openai_error)?;
resolved_api_model(
requested,
metadata.id,
AgentProvider::OpenAi,
metadata.context_window_tokens,
metadata.max_input_tokens,
)
};
self.model_cache
.lock()
.map_err(|_| {
Error::internal("model_cache_unavailable", "The model cache is unavailable.")
})?
.insert(requested.to_owned(), resolved.clone());
Ok(resolved)
}
}
impl UserIntelligence {
pub fn user_id(&self) -> &str {
&self.user_id
}
pub async fn start_agent_turn(
&self,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
mut request: kcode_codex_runtime_v2::AgentRequest,
) -> Result<AgentTurn> {
let resolved = self.service.resolve_agent_model(&request.model).await?;
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 mut operation = self
.service
.active_operations
.register_request(operation_id, parent_operation_id)?;
let requested_model = resolved.requested_model.clone();
request.model = resolved.provider_model.clone();
let (inner, actual_model, provider_request_id) = match resolved.provider {
AgentProvider::Codex => {
let inner = self.account_result(
operation_id,
parent_operation_id,
"agent_turn",
&requested_model,
self.service
.agent
.start_turn(request)
.await
.map_err(codex_v2_error),
)?;
(
AgentTurnBackend::Codex(inner),
resolved.provider_model,
None,
)
}
AgentProvider::OpenAi => {
let openai = self.service.openai.as_ref().ok_or_else(|| {
Error::unavailable("provider_not_configured", "OpenAI is not configured.")
})?;
let provider_request = openai_agent_request(&request);
let provider_input = provider_request_json(&request);
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(request.timeout, openai.agent_turn(provider_request)) => {
result
.map_err(|_| Error::provider("provider_timeout", "OpenAI agent turn timed out."))
.and_then(|result| result.map_err(openai_error))
}
};
let result = self.account_result(
operation_id,
parent_operation_id,
"agent_turn",
&requested_model,
result,
)?;
let actual_model = result.model.clone();
let provider_request_id = Some(result.response_id.clone());
(
AgentTurnBackend::Buffered(buffered_openai_turn(provider_input, result)),
actual_model,
provider_request_id,
)
}
AgentProvider::Gemini => {
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable("provider_not_configured", "Gemini is not configured.")
})?;
let provider_request = gemini_agent_request(&request);
let provider_input = provider_request_json(&request);
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(request.timeout, gemini.agent_turn(provider_request)) => {
result
.map_err(|_| Error::provider("provider_timeout", "Gemini agent turn timed out."))
.and_then(|result| result.map_err(gemini_error))
}
};
let result = self.account_result(
operation_id,
parent_operation_id,
"agent_turn",
&requested_model,
result,
)?;
let actual_model = result.model.clone();
let provider_request_id = Some(result.interaction_id.clone());
(
AgentTurnBackend::Buffered(buffered_gemini_turn(provider_input, result)),
actual_model,
provider_request_id,
)
}
};
Ok(AgentTurn {
inner,
operation,
user_id: self.user_id.clone(),
requested_model,
actual_model,
provider_request_id,
operation_id,
parent_operation_id,
receipts: self.service.receipts.clone(),
receipt: None,
})
}
pub async fn search(&self, request: SearchRequest) -> Result<Accounted<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, receipt) = 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(
request.operation_id,
request.parent_operation_id,
"web_search",
&request.model,
result,
)?;
let interaction = result.interaction;
let usage = gemini_usage(&interaction.usage);
let cost = gemini_cost(&interaction.cost);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
"web_search",
&request.model,
&interaction.model,
Some(usage),
Some(interaction.id.clone()),
None,
Some(cost.clone()),
)?;
let answer = interaction
.text
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| {
Error::provider("provider_error", "Gemini search returned no answer text.")
})
.map_err(|error| error.with_receipt(receipt.clone()))?;
(
SearchResponse {
answer,
sources: result
.sources
.into_iter()
.map(|source| WebSource {
title: source.title,
url: source.url,
})
.collect(),
model: interaction.model,
usage: Some(usage),
cost: Some(cost),
},
receipt,
)
} 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(
request.operation_id,
request.parent_operation_id,
"web_search",
&request.model,
result,
)?;
let usage = result.usage.as_ref().map(codex_usage);
let cost = codex_search_cost(&request.model, usage);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
"web_search",
&request.model,
&request.model,
usage,
None,
None,
cost.clone(),
)?;
(
SearchResponse {
answer: result.answer,
sources: result
.sources
.into_iter()
.map(|source| WebSource {
title: source.title,
url: source.url,
})
.collect(),
model: request.model.clone(),
usage,
cost,
},
receipt,
)
};
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(Accounted {
value: response,
receipt,
})
}
pub async fn analyze_audio(
&self,
request: AudioAnalysisRequest,
) -> Result<Accounted<AudioAnalysisResponse>> {
validate_operation(&request.operation)?;
validate_prompt(&request.prompt)?;
validate_model(&request.model)?;
if request.media.kind != MediaKind::Audio {
return Err(Error::invalid("audio analysis 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("audio analysis requires a supported exact Gemini model")
})?;
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"Gemini audio analysis is not configured.",
)
})?;
let media = gemini_media(&request.media)?;
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 = request
.schema
.map(StructuredOutput::new)
.transpose()
.map_err(gemini_error)?;
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 audio analysis timed out."))
.and_then(|result| result.map_err(gemini_error)),
};
let result = self.account_result(
request.operation_id,
request.parent_operation_id,
&request.operation,
&request.model,
result,
)?;
let usage = gemini_usage(&result.usage);
let cost = gemini_cost(&result.cost);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
&request.operation,
&request.model,
&result.model,
Some(usage),
Some(result.id.clone()),
None,
Some(cost.clone()),
)?;
if result.status != GeminiCompletionStatus::Completed {
return Err(Error::provider(
"provider_incomplete",
"Gemini audio analysis did not complete.",
)
.with_receipt(receipt));
}
let text = result
.text
.filter(|text| !text.trim().is_empty())
.ok_or_else(|| {
Error::provider(
"provider_empty_output",
"Gemini audio analysis returned no text.",
)
})
.map_err(|error| error.with_receipt(receipt.clone()))?;
Ok(Accounted {
value: AudioAnalysisResponse {
model: result.model,
text,
usage,
cost,
},
receipt,
})
}
pub async fn transcribe_structured_audio(
&self,
request: StructuredAudioRequest,
) -> Result<Accounted<StructuredAudioResponse>> {
self.analyze_audio(AudioAnalysisRequest {
operation: request.operation,
prompt: request.prompt,
model: request.model,
media: request.media,
schema: Some(request.schema),
max_output_tokens: request.max_output_tokens,
operation_id: request.operation_id,
parent_operation_id: request.parent_operation_id,
})
.await
.map(|response| {
response.map(|response| StructuredAudioResponse {
model: response.model,
text: response.text,
usage: response.usage,
cost: response.cost,
})
})
}
pub async fn generate_text(
&self,
request: TextGenerationRequest,
) -> Result<Accounted<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(3 * 60 * 60) {
return Err(Error::invalid(
"generation timeout must be between 1 second and 3 hours",
));
}
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_id,
request.parent_operation_id,
&request.operation,
&request.model,
result,
)?;
let usage = result.usage.as_ref().map(codex_usage);
let cost = usage.and_then(|usage| estimate_token_cost(&request.model, usage));
let receipt = self.record_tokens(
request.operation_id,
request.parent_operation_id,
&request.operation,
&request.model,
&request.model,
usage,
None,
Some(result.thread_id.clone()),
)?;
Ok(Accounted {
value: TextGenerationResponse {
model: request.model,
text: result.answer,
thread_id: result.thread_id,
usage,
cost,
},
receipt,
})
}
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<Accounted<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, receipt) = 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(
request.operation_id,
request.parent_operation_id,
"transcribe_audio",
&request.model,
result,
)?;
let metering = result
.usage
.map(transcription_usage)
.unwrap_or(Metering::Unavailable);
let cost = estimate_cost(&request.model, &metering);
let receipt = self.record_metering(
request.operation_id,
request.parent_operation_id,
"transcribe_audio",
&request.model,
&request.model,
metering.clone(),
None,
None,
)?;
(
TranscriptionResponse {
model: request.model.clone(),
text: result.text,
metering,
cost,
},
receipt,
)
} 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(
request.operation_id,
request.parent_operation_id,
"transcribe_audio",
&request.model,
result,
)?;
let metering = Metering::Tokens(gemini_usage(&result.usage));
let cost = gemini_cost(&result.cost);
let receipt = self.record_metering_with_cost(
request.operation_id,
request.parent_operation_id,
"transcribe_audio",
&request.model,
&result.model,
metering.clone(),
Some(result.id.clone()),
None,
Some(cost.clone()),
)?;
let text = result
.text
.filter(|text| !text.trim().is_empty())
.ok_or_else(|| {
Error::provider(
"empty_transcription",
"Gemini returned no transcription text.",
)
})
.map_err(|error| error.with_receipt(receipt.clone()))?;
(
TranscriptionResponse {
model: result.model,
text,
metering,
cost: Some(cost),
},
receipt,
)
};
Ok(Accounted {
value: response,
receipt,
})
}
pub async fn annotate(
&self,
request: AnnotationRequest,
) -> Result<Accounted<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, receipt) = 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(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
result,
)?;
let usage = gemini_usage(&result.usage);
let cost = gemini_cost(&result.cost);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
&result.model,
Some(usage),
Some(result.id.clone()),
None,
Some(cost.clone()),
)?;
let text = result
.text
.filter(|text| !text.trim().is_empty())
.ok_or_else(|| {
Error::provider("empty_annotation", "Gemini returned no annotation text.")
})
.map_err(|error| error.with_receipt(receipt.clone()))?;
(
AnnotationResponse {
complete: result.status == GeminiCompletionStatus::Completed,
model: result.model,
file_name,
content_type,
text,
incomplete_reason: None,
usage: Some(usage),
cost: Some(cost),
},
receipt,
)
} 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(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
result,
)?;
let usage = result.usage.as_ref().map(openai_image_usage);
let cost = result.usage.as_ref().map(openai_image_analysis_cost);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
&result.model,
usage,
None,
None,
cost.clone(),
)?;
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,
cost,
},
receipt,
)
} 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(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
turn_result,
)?;
let completed = loop {
let event = tokio::select! {
_ = operation.cancelled() => {
turn.cancel();
let receipt = self.record_metering(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(Error::cancelled().with_receipt(receipt));
}
event = turn.next_event() => event,
};
match event {
Some(Ok(kcode_codex_runtime_v2::AgentEvent::ProviderInput(_))) => {}
Some(Ok(kcode_codex_runtime_v2::AgentEvent::UsageUpdated(_))) => {}
Some(Ok(kcode_codex_runtime_v2::AgentEvent::ToolCall(_))) => {
turn.cancel();
let receipt = self.record_metering(
request.operation_id,
request.parent_operation_id,
"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.",
)
.with_receipt(receipt));
}
Some(Ok(kcode_codex_runtime_v2::AgentEvent::Completed(completed))) => {
break completed;
}
Some(Err(error)) => {
let receipt = self.record_metering(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(codex_v2_error(error).with_receipt(receipt));
}
None => {
let receipt = self.record_metering(
request.operation_id,
request.parent_operation_id,
"annotate_media",
&request.model,
&request.model,
Metering::Unavailable,
None,
None,
)?;
return Err(Error::provider(
"empty_annotation",
"Codex ended without annotation text.",
)
.with_receipt(receipt));
}
}
};
let usage = completed.usage.as_ref().map(codex_v2_usage);
let cost = usage.and_then(|usage| estimate_token_cost(&request.model, usage));
let receipt = self.record_tokens(
request.operation_id,
request.parent_operation_id,
"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,
cost,
},
receipt,
)
} else {
return Err(Error::invalid(format!(
"unsupported exact annotation model {}",
request.model
)));
};
Ok(Accounted {
value: response,
receipt,
})
}
pub async fn generate_image(&self, request: ImageRequest) -> Result<Accounted<ImageResponse>> {
validate_model(&request.model)?;
if request.prompt.trim().is_empty()
|| request.prompt.chars().count() > MAX_IMAGE_PROMPT_CHARACTERS
{
return Err(Error::invalid(format!(
"image prompt must contain 1 through {MAX_IMAGE_PROMPT_CHARACTERS} characters"
)));
}
if request
.references
.iter()
.any(|media| media.kind != MediaKind::Image)
{
return Err(Error::invalid("image references must all be images"));
}
let mut operation = self
.service
.active_operations
.register_request(request.operation_id, request.parent_operation_id)?;
let operation_name = if request.references.is_empty() {
"generate_image"
} else {
"edit_image"
};
if request.model == kcode_openai_api::GPT_IMAGE_2 {
let openai = self.service.openai.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"OpenAI image generation is not configured.",
)
})?;
let result = if request.references.is_empty() {
let provider_request = OpenAiImageRequest::new(request.prompt);
tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
IMAGE_OPERATION_TIMEOUT,
openai.generate_image(provider_request),
) => result
.map_err(|_| Error::provider("provider_timeout", "OpenAI image generation timed out."))
.and_then(|result| result.map_err(openai_error)),
}
} else {
let mut images = request
.references
.into_iter()
.map(|media| {
OpenAiImageInput::new(
openai_image_media_type(&media.content_type)?,
media.bytes,
)
.map_err(openai_error)
})
.collect::<Result<Vec<_>>>()?;
let first = images.remove(0);
let mut provider_request = OpenAiImageEditRequest::new(first, request.prompt);
provider_request.images.extend(images);
tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
IMAGE_OPERATION_TIMEOUT,
openai.edit_image(provider_request),
) => result
.map_err(|_| Error::provider("provider_timeout", "OpenAI image editing timed out."))
.and_then(|result| result.map_err(openai_error)),
}
};
let result = self.account_result(
request.operation_id,
request.parent_operation_id,
operation_name,
&request.model,
result,
)?;
let usage = result.usage.as_ref().map(openai_generation_usage);
let cost = result.usage.as_ref().map(openai_generation_cost);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
operation_name,
&request.model,
kcode_openai_api::GPT_IMAGE_2,
usage,
result.request_id,
None,
cost.clone(),
)?;
Ok(Accounted {
value: ImageResponse {
model: kcode_openai_api::GPT_IMAGE_2.into(),
content_type: result.image.format.mime_type().into(),
bytes: result.image.data,
usage,
cost,
},
receipt,
})
} else if request.model == kcode_gemini_api::NANO_BANANA_PRO {
let gemini = self.service.gemini.as_ref().ok_or_else(|| {
Error::unavailable(
"provider_not_configured",
"Gemini image generation is not configured.",
)
})?;
let mut provider_request = NanoBananaProRequest::new(request.prompt);
provider_request.images = request
.references
.into_iter()
.map(|media| {
GeminiMediaInput::image(&media.content_type, media.bytes).map_err(gemini_error)
})
.collect::<Result<Vec<_>>>()?;
let result = tokio::select! {
_ = operation.cancelled() => Err(Error::cancelled()),
result = tokio::time::timeout(
IMAGE_OPERATION_TIMEOUT,
gemini.nano_banana_pro(provider_request),
) => result
.map_err(|_| Error::provider("provider_timeout", "Gemini image generation timed out."))
.and_then(|result| result.map_err(gemini_error)),
};
let result = self.account_result(
request.operation_id,
request.parent_operation_id,
operation_name,
&request.model,
result,
)?;
let usage = gemini_usage(&result.usage);
let cost = gemini_cost(&result.cost);
let receipt = self.record_tokens_with_cost(
request.operation_id,
request.parent_operation_id,
operation_name,
&request.model,
&result.model,
Some(usage),
Some(result.id.clone()),
None,
Some(cost.clone()),
)?;
let mut images = result.images;
if images.len() != 1 {
return Err(Error::provider(
"provider_error",
"Gemini did not return exactly one generated image.",
)
.with_receipt(receipt));
}
let image = images.remove(0);
Ok(Accounted {
value: ImageResponse {
model: result.model,
content_type: image.mime_type,
bytes: image.data,
usage: Some(usage),
cost: Some(cost),
},
receipt,
})
} else {
Err(Error::invalid(format!(
"unsupported exact image model {}",
request.model
)))
}
}
#[allow(clippy::too_many_arguments)]
fn record_tokens(
&self,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
operation: &str,
requested_model: &str,
actual_model: &str,
usage: Option<TokenUsage>,
provider_request_id: Option<String>,
provider_thread_id: Option<String>,
) -> Result<UsageReceipt> {
self.record_metering_with_cost(
operation_id,
parent_operation_id,
operation,
requested_model,
actual_model,
usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
provider_request_id,
provider_thread_id,
None,
)
}
#[allow(clippy::too_many_arguments)]
fn record_tokens_with_cost(
&self,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
operation: &str,
requested_model: &str,
actual_model: &str,
usage: Option<TokenUsage>,
provider_request_id: Option<String>,
provider_thread_id: Option<String>,
cost: Option<CostEstimate>,
) -> Result<UsageReceipt> {
self.record_metering_with_cost(
operation_id,
parent_operation_id,
operation,
requested_model,
actual_model,
usage.map(Metering::Tokens).unwrap_or(Metering::Unavailable),
provider_request_id,
provider_thread_id,
cost,
)
}
#[allow(clippy::too_many_arguments)]
fn record_metering(
&self,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
operation: &str,
requested_model: &str,
actual_model: &str,
metering: Metering,
provider_request_id: Option<String>,
provider_thread_id: Option<String>,
) -> Result<UsageReceipt> {
self.record_metering_with_cost(
operation_id,
parent_operation_id,
operation,
requested_model,
actual_model,
metering,
provider_request_id,
provider_thread_id,
None,
)
}
#[allow(clippy::too_many_arguments)]
fn record_metering_with_cost(
&self,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
operation: &str,
requested_model: &str,
actual_model: &str,
metering: Metering,
provider_request_id: Option<String>,
provider_thread_id: Option<String>,
cost: Option<CostEstimate>,
) -> Result<UsageReceipt> {
let mut receipt = UsageReceipt::new(
self.user_id.clone(),
operation_id,
parent_operation_id,
operation,
requested_model,
actual_model,
metering,
);
receipt.provider_request_id = provider_request_id;
receipt.provider_thread_id = provider_thread_id;
if cost.is_some() {
receipt.cost = cost;
}
self.service.receipts.record(&receipt)?;
Ok(receipt)
}
fn account_result<T>(
&self,
operation_id: Uuid,
parent_operation_id: Option<Uuid>,
operation: &str,
requested_model: &str,
result: Result<T>,
) -> Result<T> {
match result {
Ok(value) => Ok(value),
Err(error) => {
let receipt = self.record_metering(
operation_id,
parent_operation_id,
operation,
requested_model,
requested_model,
Metering::Unavailable,
None,
None,
)?;
Err(error.with_receipt(receipt))
}
}
}
}
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 fn receipt(&self) -> Option<&UsageReceipt> {
self.receipt.as_ref()
}
pub fn finish_unavailable(&mut self) -> Result<&UsageReceipt> {
if self.receipt.is_none() {
self.record_unavailable()?;
}
Ok(self.receipt.as_ref().expect("receipt was just recorded"))
}
pub async fn next_event(&mut self) -> Result<Option<kcode_codex_runtime_v2::AgentEvent>> {
let event = match &mut self.inner {
AgentTurnBackend::Codex(inner) => tokio::select! {
_ = self.operation.cancelled() => {
inner.cancel();
let receipt = self.record_unavailable()?;
return Err(Error::cancelled().with_receipt(receipt));
}
event = inner.next_event() => match event {
Some(Ok(event)) => Some(event),
Some(Err(error)) => {
let receipt = self.record_unavailable()?;
return Err(codex_v2_error(error).with_receipt(receipt));
}
None => {
self.record_unavailable()?;
None
},
}
},
AgentTurnBackend::Buffered(inner) => {
if *self.operation.cancellation.borrow() {
let receipt = self.record_unavailable()?;
return Err(Error::cancelled().with_receipt(receipt));
}
inner.events.pop_front()
}
};
if let Some(kcode_codex_runtime_v2::AgentEvent::Completed(completed)) = &event
&& self.receipt.is_none()
{
let mut receipt = UsageReceipt::new(
self.user_id.clone(),
self.operation_id,
self.parent_operation_id,
"agent_turn",
self.requested_model.clone(),
self.actual_model.clone(),
completed
.usage
.as_ref()
.map(codex_v2_usage)
.map(Metering::Tokens)
.unwrap_or(Metering::Unavailable),
);
receipt.provider_request_id = self
.provider_request_id
.clone()
.or_else(|| Some(completed.turn_id.clone()));
receipt.provider_thread_id = Some(completed.thread_id.clone());
self.receipts.record(&receipt)?;
self.receipt = Some(receipt);
}
Ok(event)
}
fn record_unavailable(&mut self) -> Result<UsageReceipt> {
if let Some(receipt) = &self.receipt {
return Ok(receipt.clone());
}
let receipt = UsageReceipt::new(
self.user_id.clone(),
self.operation_id,
self.parent_operation_id,
"agent_turn",
self.requested_model.clone(),
self.actual_model.clone(),
Metering::Unavailable,
);
self.receipts.record(&receipt)?;
self.receipt = Some(receipt.clone());
Ok(receipt)
}
pub async fn respond(
&mut self,
call_id: &str,
result: kcode_codex_runtime_v2::ToolResult,
) -> Result<()> {
match &mut self.inner {
AgentTurnBackend::Codex(inner) => {
inner.respond(call_id, result).await.map_err(codex_v2_error)
}
AgentTurnBackend::Buffered(inner) => {
if inner.pending_call_id.as_deref() != Some(call_id) {
return Err(Error::invalid(
"tool result does not match the pending provider call",
));
}
inner.pending_call_id = None;
if let Some(mut completed) = inner.completed.take() {
completed.answer.clear();
inner
.events
.push_back(kcode_codex_runtime_v2::AgentEvent::Completed(completed));
}
Ok(())
}
}
}
}
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,
parent_cancellation: None,
})
}
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 parent_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.",
)
})?;
let mut operation = self.register(request_id)?;
operation.parent_cancellation = Some(parent_cancellation);
Ok(operation)
} 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 let Some(parent) = &mut self.parent_cancellation {
tokio::select! {
_ = cancellation_requested(&mut self.cancellation) => {}
_ = cancellation_requested(parent) => {}
}
} else {
cancellation_requested(&mut self.cancellation).await;
}
}
}
async fn cancellation_requested(cancellation: &mut watch::Receiver<bool>) {
if *cancellation.borrow() {
return;
}
while cancellation.changed().await.is_ok() {
if *cancellation.borrow() {
return;
}
}
}
impl Drop for ActiveOperation {
fn drop(&mut self) {
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_agent_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'_' | b'/' | b':')
})
|| model.starts_with("codex:")
{
return Err(Error::invalid(
"model must be an exact safe provider 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 resolved_api_model(
requested: &str,
provider_model: String,
provider: AgentProvider,
context_window_tokens: Option<u64>,
max_input_tokens: Option<u64>,
) -> ResolvedAgentModel {
let known = match provider {
AgentProvider::OpenAi if requested == "gpt-5.6" => Some((1_000_000, 700_000)),
AgentProvider::Gemini
if matches!(
requested,
"gemini-2.5-flash" | "gemini-3.1-flash-lite" | "gemini-3.1-pro-preview"
) =>
{
Some((1_000_000, 700_000))
}
_ => None,
};
let context_window_tokens = context_window_tokens
.or_else(|| known.map(|limits| limits.0))
.unwrap_or(128_000);
let max_input_tokens = max_input_tokens
.or_else(|| known.map(|limits| limits.1))
.unwrap_or_else(|| context_window_tokens.saturating_mul(70) / 100)
.min(context_window_tokens);
ResolvedAgentModel {
requested_model: requested.to_owned(),
provider_model,
provider,
context_window_tokens,
max_input_tokens,
}
}
fn openai_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> OpenAiAgentRequest {
OpenAiAgentRequest {
model: request.model.clone(),
input: request.input.clone(),
reasoning_effort: request.reasoning_effort.as_str().into(),
tools: request
.tools
.iter()
.map(|tool| OpenAiAgentTool {
name: tool.name.clone(),
description: tool.description.clone(),
input_schema: tool.input_schema.clone(),
})
.collect(),
}
}
fn gemini_agent_request(request: &kcode_codex_runtime_v2::AgentRequest) -> GeminiAgentRequest {
GeminiAgentRequest {
model: request.model.clone(),
input: request.input.clone(),
reasoning_effort: request.reasoning_effort.as_str().into(),
tools: request
.tools
.iter()
.map(|tool| GeminiAgentTool {
name: tool.name.clone(),
description: tool.description.clone(),
input_schema: tool.input_schema.clone(),
})
.collect(),
}
}
fn provider_request_json(request: &kcode_codex_runtime_v2::AgentRequest) -> String {
format!(
"{}\n",
serde_json::to_string(&json!({
"model": request.model,
"input": request.input,
"reasoningEffort": request.reasoning_effort.as_str(),
"tools": request.tools.iter().map(|tool| json!({
"name": tool.name,
"description": tool.description,
"parameters": tool.input_schema
})).collect::<Vec<_>>()
}))
.expect("provider request values always serialize")
)
}
fn buffered_openai_turn(
provider_input: String,
response: kcode_openai_api::AgentTurnResponse,
) -> BufferedAgentTurn {
let usage = response
.usage
.as_ref()
.map(|usage| kcode_codex_runtime_v2::TokenUsage {
input_tokens: usage.input_tokens,
output_tokens: usage.output_tokens,
cached_input_tokens: usage.cached_input_tokens,
reasoning_output_tokens: usage.reasoning_output_tokens,
last_input_tokens: None,
last_output_tokens: None,
});
let completed = kcode_codex_runtime_v2::CompletedTurn {
thread_id: response.response_id,
turn_id: Uuid::new_v4().to_string(),
answer: response.text,
usage,
};
buffered_turn(
provider_input,
response
.tool_call
.map(|call| kcode_codex_runtime_v2::DynamicToolCall {
call_id: call.call_id,
tool: call.name,
arguments: call.arguments,
}),
completed,
)
}
fn buffered_gemini_turn(
provider_input: String,
response: kcode_gemini_api::AgentTurnResponse,
) -> BufferedAgentTurn {
let usage = kcode_codex_runtime_v2::TokenUsage {
input_tokens: response.usage.input_tokens,
output_tokens: response
.usage
.output_tokens
.saturating_add(response.usage.thought_tokens),
cached_input_tokens: response.usage.cached_tokens,
reasoning_output_tokens: response.usage.thought_tokens,
last_input_tokens: None,
last_output_tokens: None,
};
let completed = kcode_codex_runtime_v2::CompletedTurn {
thread_id: response.interaction_id,
turn_id: Uuid::new_v4().to_string(),
answer: response.text,
usage: Some(usage),
};
buffered_turn(
provider_input,
response
.tool_call
.map(|call| kcode_codex_runtime_v2::DynamicToolCall {
call_id: call.call_id,
tool: call.name,
arguments: call.arguments,
}),
completed,
)
}
fn buffered_turn(
provider_input: String,
tool_call: Option<kcode_codex_runtime_v2::DynamicToolCall>,
completed: kcode_codex_runtime_v2::CompletedTurn,
) -> BufferedAgentTurn {
let mut events = VecDeque::from([kcode_codex_runtime_v2::AgentEvent::ProviderInput(
provider_input,
)]);
if let Some(usage) = completed.usage.clone() {
events.push_back(kcode_codex_runtime_v2::AgentEvent::UsageUpdated(usage));
}
let pending_call_id = tool_call.as_ref().map(|call| call.call_id.clone());
if let Some(call) = tool_call {
events.push_back(kcode_codex_runtime_v2::AgentEvent::ToolCall(call));
} else {
events.push_back(kcode_codex_runtime_v2::AgentEvent::Completed(
completed.clone(),
));
}
BufferedAgentTurn {
events,
pending_call_id,
completed: Some(completed),
}
}
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 gemini_cost(cost: &kcode_gemini_api::CostBreakdown) -> CostEstimate {
let accuracy = match cost.accuracy {
kcode_gemini_api::CostAccuracy::Exact => CostAccuracy::Exact,
kcode_gemini_api::CostAccuracy::Estimated => CostAccuracy::Estimated,
kcode_gemini_api::CostAccuracy::Conservative => CostAccuracy::Conservative,
};
CostEstimate {
usd_nanos: cost.total.usd_nanos(),
accuracy,
pricing_version: cost.pricing_version.clone(),
}
}
fn openai_image_analysis_cost(usage: &ImageAnalysisUsage) -> CostEstimate {
let cached = usage.cached_input_tokens.unwrap_or(0);
let cache_write = usage.cache_write_input_tokens.unwrap_or(0);
let non_cached = usage
.input_tokens
.saturating_sub(cached)
.saturating_sub(cache_write);
let long = usage.input_tokens > 272_000;
let (input_rate, cached_rate, cache_write_rate, output_rate) = if long {
(10_000, 1_000, 12_500, 45_000)
} else {
(5_000, 500, 6_250, 30_000)
};
CostEstimate {
usd_nanos: non_cached
.saturating_mul(input_rate)
.saturating_add(cached.saturating_mul(cached_rate))
.saturating_add(cache_write.saturating_mul(cache_write_rate))
.saturating_add(usage.output_tokens.saturating_mul(output_rate)),
accuracy: CostAccuracy::Exact,
pricing_version: PRICING_VERSION.into(),
}
}
fn openai_generation_cost(usage: &OpenAiGenerationUsage) -> CostEstimate {
let detailed_input = usage
.input_details
.text_tokens
.saturating_add(usage.input_details.image_tokens);
let output_details = usage.output_details.as_ref();
let detailed_output = output_details
.map(|details| details.text_tokens.saturating_add(details.image_tokens))
.unwrap_or(0);
let input = usage
.input_details
.text_tokens
.saturating_mul(5_000)
.saturating_add(usage.input_details.image_tokens.saturating_mul(8_000))
.saturating_add(
usage
.input_tokens
.saturating_sub(detailed_input)
.saturating_mul(8_000),
);
let output = detailed_output
.saturating_add(usage.output_tokens.saturating_sub(detailed_output))
.saturating_mul(30_000);
CostEstimate {
usd_nanos: input.saturating_add(output),
accuracy: if detailed_input == usage.input_tokens
&& output_details.is_some()
&& detailed_output == usage.output_tokens
{
CostAccuracy::Exact
} else {
CostAccuracy::Estimated
},
pricing_version: PRICING_VERSION.into(),
}
}
fn codex_search_cost(model: &str, usage: Option<TokenUsage>) -> Option<CostEstimate> {
let metering = Metering::Tokens(usage?);
estimate_cost(model, &metering).map(|cost| {
cost.with_surcharge(10_000_000, CostAccuracy::Estimated)
})
}
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 openai_generation_usage(usage: &OpenAiGenerationUsage) -> TokenUsage {
TokenUsage {
input_tokens: usage.input_tokens,
cached_input_tokens: 0,
thinking_tokens: 0,
output_tokens: usage.output_tokens,
}
}
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());
}
#[tokio::test]
async fn child_operations_have_independent_ids_and_inherit_parent_cancellation() {
let operations = ActiveOperations::default();
let parent_id = Uuid::new_v4();
let child_id = Uuid::new_v4();
let _parent = operations.register(parent_id).unwrap();
let mut child = operations
.register_request(child_id, Some(parent_id))
.unwrap();
assert!(operations.cancel(child_id).unwrap());
tokio::time::timeout(Duration::from_millis(50), child.cancelled())
.await
.unwrap();
assert!(operations.cancel(parent_id).unwrap());
let child_id = Uuid::new_v4();
let mut inherited = operations
.register_request(child_id, Some(parent_id))
.unwrap();
tokio::time::timeout(Duration::from_millis(50), inherited.cancelled())
.await
.unwrap();
}
}