pub use codex_api::ResponseEvent; use codex_config::types::Personality; use codex_protocol::error::Result; use codex_protocol::models::BaseInstructions; use codex_protocol::models::ContentItem; use codex_protocol::models::FunctionCallOutputContentItem; use codex_protocol::models::ResponseItem; use codex_protocol::protocol::InterAgentCommunication; use codex_tools::ToolSpec; use futures::Stream; use serde_json::Value; use std::pin::Pin; use std::task::Context; use std::task::Poll; use tokio::sync::mpsc; use tokio_util::sync::CancellationToken; /// API request payload for a single model turn #[derive(Debug, Clone)] pub struct Prompt { /// Conversation context input items. pub input: Vec, /// Tools available to the model, including additional tools sourced from /// external MCP servers. pub(crate) tools: Vec, /// Whether parallel tool calls are permitted for this prompt. pub(crate) parallel_tool_calls: bool, pub base_instructions: BaseInstructions, /// Optionally specify the personality of the model. pub personality: Option, /// Optional the output schema for the model's response. pub output_schema: Option, /// Whether the Responses API should strictly validate `output_schema`. pub output_schema_strict: bool, } impl Default for Prompt { fn default() -> Self { Self { input: Vec::new(), tools: Vec::new(), parallel_tool_calls: false, base_instructions: BaseInstructions::default(), personality: None, output_schema: None, output_schema_strict: true, } } } impl Prompt { pub(crate) fn get_formatted_input(&self) -> Vec { self.input .iter() .cloned() .map(|item| { let ResponseItem::Message { role, content, .. } = &item else { return item; }; if role != "assistant" { return item; } InterAgentCommunication::from_message_content(content) .filter(|communication| communication.encrypted_content.is_some()) .map(|communication| communication.to_model_input_item()) .unwrap_or(item) }) .collect() } pub(crate) fn get_formatted_input_for_request( &self, use_responses_lite: bool, ) -> Vec { let mut input = self.get_formatted_input(); if use_responses_lite { strip_image_details(&mut input); } input } } fn strip_image_details(items: &mut [ResponseItem]) { for item in items { match item { ResponseItem::Message { content, .. } => { for content_item in content { if let ContentItem::InputImage { detail, .. } = content_item { *detail = None; } } } ResponseItem::FunctionCallOutput { output, .. } | ResponseItem::CustomToolCallOutput { output, .. } => { if let Some(content) = output.content_items_mut() { for content_item in content { if let FunctionCallOutputContentItem::InputImage { detail, .. } = content_item { *detail = None; } } } } ResponseItem::Reasoning { .. } | ResponseItem::AgentMessage { .. } | ResponseItem::LocalShellCall { .. } | ResponseItem::FunctionCall { .. } | ResponseItem::ToolSearchCall { .. } | ResponseItem::CustomToolCall { .. } | ResponseItem::ToolSearchOutput { .. } | ResponseItem::WebSearchCall { .. } | ResponseItem::ImageGenerationCall { .. } | ResponseItem::Compaction { .. } | ResponseItem::CompactionTrigger | ResponseItem::ContextCompaction { .. } | ResponseItem::Other => {} } } } pub struct ResponseStream { pub(crate) rx_event: mpsc::Receiver>, /// Signals the mapper task that the consumer stopped polling before the /// provider stream reached its own terminal event. pub(crate) consumer_dropped: CancellationToken, } impl Stream for ResponseStream { type Item = Result; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { self.rx_event.poll_recv(cx) } } impl Drop for ResponseStream { fn drop(&mut self) { self.consumer_dropped.cancel(); } } #[cfg(test)] #[path = "client_common_tests.rs"] mod tests;