use std::{collections::HashSet, path::PathBuf, process::Stdio, sync::Arc, time::Duration};
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
use serde_json::{Value, json};
use tokio::{
io::{AsyncBufReadExt, AsyncWrite, AsyncWriteExt, BufReader},
process::{Child, Command},
sync::{mpsc, watch},
};
use crate::{
AgentEvent, AgentRequest, CompletedTurn, DynamicToolCall, Error, ErrorKind, ImageInput,
ImageTurnRequest, Result, TokenUsage, ToolResult,
};
pub const DEFAULT_CODEX_EXECUTABLE: &str = "codex-safe";
const TOOL_OUTPUT_TOKEN_LIMIT: usize = 180_000;
const MAX_INPUT_CHARACTERS: usize = 1_048_576;
const MAX_IMAGE_COUNT: usize = 8;
const MAX_IMAGE_BYTES: usize = 20 * 1024 * 1024;
const STARTUP_TIMEOUT: Duration = Duration::from_secs(90);
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct CodexConfig {
pub executable: String,
pub working_directory: PathBuf,
pub base_instruction: String,
pub model_catalog: Option<PathBuf>,
}
impl Default for CodexConfig {
fn default() -> Self {
Self {
executable: DEFAULT_CODEX_EXECUTABLE.into(),
working_directory: std::env::temp_dir(),
base_instruction: String::new(),
model_catalog: None,
}
}
}
#[derive(Clone, Debug)]
pub struct Codex {
config: Arc<CodexConfig>,
}
impl Codex {
pub async fn open(config: CodexConfig) -> Result<Self> {
validate_config(&config)?;
validate_chatgpt_login(&config.executable).await?;
Ok(Self {
config: Arc::new(config),
})
}
pub async fn start_turn(&self, request: AgentRequest) -> Result<AgentTurn> {
validate_request(&request)?;
self.start_validated_turn(request, Vec::new())
}
pub async fn start_image_turn(&self, request: ImageTurnRequest) -> Result<AgentTurn> {
let ImageTurnRequest {
prompt,
model,
images,
reasoning_effort,
timeout,
} = request;
validate_images(&images)?;
let request = AgentRequest {
input: prompt,
model,
reasoning_effort,
previous_thread_id: None,
tools: Vec::new(),
ephemeral: true,
timeout,
};
validate_request(&request)?;
self.start_validated_turn(request, images)
}
fn start_validated_turn(
&self,
request: AgentRequest,
images: Vec<ImageInput>,
) -> Result<AgentTurn> {
let child = spawn_app_server(&self.config)?;
let (event_sender, events) = mpsc::channel(32);
let (tool_results, result_receiver) = mpsc::channel(8);
let (cancel, cancelled) = watch::channel(false);
let config = self.config.clone();
let timeout = request.timeout;
tokio::spawn(async move {
let result = tokio::select! {
_ = cancellation(cancelled) => Err(Error::new(ErrorKind::Cancelled, "Codex turn was cancelled")),
result = tokio::time::timeout(timeout, run_protocol(child, &config, request, images, &event_sender, result_receiver)) => {
match result {
Ok(result) => result,
Err(_) => Err(Error::new(ErrorKind::Timeout, "Codex turn timed out")),
}
}
};
if let Err(error) = result {
let _ = event_sender.send(Err(error)).await;
}
});
Ok(AgentTurn {
events,
tool_results,
cancel,
})
}
}
pub struct AgentTurn {
events: mpsc::Receiver<Result<AgentEvent>>,
tool_results: mpsc::Sender<(String, ToolResult)>,
cancel: watch::Sender<bool>,
}
impl AgentTurn {
pub async fn next_event(&mut self) -> Option<Result<AgentEvent>> {
self.events.recv().await
}
pub async fn respond(&self, call_id: impl Into<String>, result: ToolResult) -> Result<()> {
self.tool_results
.send((call_id.into(), result))
.await
.map_err(|_| Error::new(ErrorKind::Cancelled, "Codex turn is no longer running"))
}
pub fn cancel(&self) {
let _ = self.cancel.send(true);
}
}
impl Drop for AgentTurn {
fn drop(&mut self) {
let _ = self.cancel.send(true);
}
}
async fn cancellation(mut cancelled: watch::Receiver<bool>) {
if *cancelled.borrow() {
return;
}
while cancelled.changed().await.is_ok() {
if *cancelled.borrow() {
return;
}
}
}
fn validate_config(config: &CodexConfig) -> Result<()> {
if config.executable.trim().is_empty() {
return Err(Error::new(
ErrorKind::InvalidInput,
"Codex executable must not be empty",
));
}
if config.base_instruction.chars().count() > MAX_INPUT_CHARACTERS {
return Err(Error::new(
ErrorKind::InvalidInput,
"Codex base instruction is too large",
));
}
Ok(())
}
fn validate_request(request: &AgentRequest) -> Result<()> {
if request.input.trim().is_empty() {
return Err(Error::new(
ErrorKind::InvalidInput,
"Codex turn input must not be empty",
));
}
if request.input.chars().count() > MAX_INPUT_CHARACTERS {
return Err(Error::new(
ErrorKind::InvalidInput,
format!("Codex turn input exceeds {MAX_INPUT_CHARACTERS} characters"),
));
}
if request.model.trim().is_empty() {
return Err(Error::new(
ErrorKind::InvalidInput,
"Codex model must not be empty",
));
}
if request.timeout.is_zero() {
return Err(Error::new(
ErrorKind::InvalidInput,
"Codex timeout must be greater than zero",
));
}
let mut names = HashSet::new();
for tool in &request.tools {
if tool.name.trim().is_empty() || tool.description.trim().is_empty() {
return Err(Error::new(
ErrorKind::InvalidInput,
"dynamic tool names and descriptions must not be empty",
));
}
if !tool.input_schema.is_object() {
return Err(Error::new(
ErrorKind::InvalidInput,
format!(
"dynamic tool '{}' must use an object JSON Schema",
tool.name
),
));
}
if !names.insert(tool.name.as_str()) {
return Err(Error::new(
ErrorKind::InvalidInput,
format!("dynamic tool '{}' is duplicated", tool.name),
));
}
}
Ok(())
}
fn validate_images(images: &[ImageInput]) -> Result<()> {
if images.is_empty() {
return Err(Error::new(
ErrorKind::InvalidInput,
"Codex image turn requires at least one image",
));
}
if images.len() > MAX_IMAGE_COUNT {
return Err(Error::new(
ErrorKind::InvalidInput,
format!("Codex image turn supports at most {MAX_IMAGE_COUNT} images"),
));
}
let total_bytes = images
.iter()
.try_fold(0usize, |total, image| {
total.checked_add(image.bytes().len())
})
.ok_or_else(|| {
Error::new(
ErrorKind::InvalidInput,
"Codex image input byte count overflowed",
)
})?;
if total_bytes > MAX_IMAGE_BYTES {
return Err(Error::new(
ErrorKind::InvalidInput,
format!("Codex image turn exceeds the {MAX_IMAGE_BYTES}-byte aggregate image limit"),
));
}
Ok(())
}
async fn validate_chatgpt_login(executable: &str) -> Result<()> {
let output = tokio::time::timeout(
STARTUP_TIMEOUT,
Command::new(executable)
.args(["login", "status"])
.env_remove("OPENAI_API_KEY")
.env_remove("CODEX_API_KEY")
.output(),
)
.await
.map_err(|_| Error::new(ErrorKind::Unavailable, "Codex login check timed out"))?
.map_err(|_| {
Error::new(
ErrorKind::Unavailable,
format!("Codex sandbox launcher '{executable}' could not be started"),
)
})?;
let status = format!(
"{}\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
if !output.status.success() || !status.to_ascii_lowercase().contains("chatgpt") {
return Err(Error::new(
ErrorKind::Authentication,
format!("'{executable}' must be logged in with ChatGPT"),
));
}
Ok(())
}
fn spawn_app_server(config: &CodexConfig) -> Result<Child> {
app_server_command(config).spawn().map_err(|_| {
Error::new(
ErrorKind::Unavailable,
format!(
"Codex sandbox launcher '{}' could not start app-server",
config.executable
),
)
})
}
fn app_server_command(config: &CodexConfig) -> Command {
let mut command = Command::new(&config.executable);
command
.arg("-c")
.arg("web_search=\"disabled\"")
.arg("-c")
.arg("mcp_servers={}")
.arg("-c")
.arg("features.shell_tool=false")
.arg("-c")
.arg("features.apps=false")
.arg("-c")
.arg("features.browser_use=false")
.arg("-c")
.arg("features.computer_use=false")
.arg("-c")
.arg("features.goals=false")
.arg("-c")
.arg("features.hooks=false")
.arg("-c")
.arg("features.image_generation=false")
.arg("-c")
.arg("features.multi_agent=false")
.arg("-c")
.arg("features.plugins=false")
.arg("-c")
.arg("features.tool_suggest=false")
.arg("-c")
.arg("features.remote_plugin=false")
.arg("-c")
.arg("model_auto_compact_token_limit=9223372036854775807")
.arg("-c")
.arg(format!("tool_output_token_limit={TOOL_OUTPUT_TOKEN_LIMIT}"));
if let Some(path) = &config.model_catalog {
command.arg("-c").arg(format!(
"model_catalog_json={}",
serde_json::to_string(path.to_string_lossy().as_ref())
.expect("serializing a path string cannot fail")
));
}
command
.arg("app-server")
.arg("--stdio")
.current_dir(&config.working_directory)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.env_remove("OPENAI_API_KEY")
.env_remove("CODEX_API_KEY")
.kill_on_drop(true);
command
}
async fn run_protocol(
mut child: Child,
config: &CodexConfig,
request: AgentRequest,
images: Vec<ImageInput>,
events: &mpsc::Sender<Result<AgentEvent>>,
mut tool_results: mpsc::Receiver<(String, ToolResult)>,
) -> Result<()> {
let stdin = child.stdin.take().ok_or_else(|| {
Error::new(
ErrorKind::Unavailable,
"Codex app-server standard input is unavailable",
)
})?;
let stdout = child.stdout.take().ok_or_else(|| {
Error::new(
ErrorKind::Unavailable,
"Codex app-server standard output is unavailable",
)
})?;
let mut writer = stdin;
let mut lines = BufReader::new(stdout).lines();
write_record(
&mut writer,
events,
json!({
"id":1,
"method":"initialize",
"params":{
"clientInfo":{"name":"kcode-codex-runtime-v2","version":env!("CARGO_PKG_VERSION")},
"capabilities":{"experimentalApi":true}
}
}),
)
.await?;
let initialized = read_response(&mut lines, 1).await?;
require_result(&initialized, "initialize")?;
write_record(
&mut writer,
events,
json!({"method":"initialized","params":{}}),
)
.await?;
let thread_method;
let thread_params;
if let Some(thread_id) = request.previous_thread_id.as_deref() {
thread_method = "thread/resume";
thread_params = json!({
"threadId":thread_id,
"model":request.model,
"cwd":config.working_directory,
"approvalPolicy":"never",
"sandbox":"read-only",
"baseInstructions":config.base_instruction,
});
} else {
thread_method = "thread/start";
let tools = request
.tools
.iter()
.map(|tool| {
json!({
"type":"function",
"name":tool.name,
"description":tool.description,
"inputSchema":tool.input_schema,
})
})
.collect::<Vec<_>>();
thread_params = json!({
"model":request.model,
"cwd":config.working_directory,
"approvalPolicy":"never",
"sandbox":"read-only",
"baseInstructions":config.base_instruction,
"developerInstructions":"",
"dynamicTools":tools,
"ephemeral":request.ephemeral,
"environments":[],
});
}
write_record(
&mut writer,
events,
json!({"id":2,"method":thread_method,"params":thread_params}),
)
.await?;
let thread_response = read_response(&mut lines, 2).await?;
let thread_result = require_result(&thread_response, thread_method)?;
let thread_id = thread_result
.pointer("/thread/id")
.and_then(Value::as_str)
.ok_or_else(|| Error::new(ErrorKind::Protocol, "Codex omitted its thread ID"))?
.to_owned();
write_record(
&mut writer,
events,
json!({
"id":3,
"method":"turn/start",
"params":{
"threadId":thread_id,
"input":turn_input(&request.input, &images),
"effort":request.reasoning_effort.as_str(),
"approvalPolicy":"never",
}
}),
)
.await?;
let turn_response = read_response(&mut lines, 3).await?;
let turn_result = require_result(&turn_response, "turn/start")?;
let turn_id = turn_result
.pointer("/turn/id")
.and_then(Value::as_str)
.ok_or_else(|| Error::new(ErrorKind::Protocol, "Codex omitted its turn ID"))?
.to_owned();
let mut answer = String::new();
let mut usage = None;
while let Some(line) = lines.next_line().await.map_err(|_| {
Error::new(
ErrorKind::Unavailable,
"Codex app-server output could not be read",
)
})? {
let message: Value = serde_json::from_str(&line)
.map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
match message.get("method").and_then(Value::as_str) {
Some("item/tool/call") => {
let id = message.get("id").cloned().ok_or_else(|| {
Error::new(
ErrorKind::Protocol,
"Codex tool request omitted its JSON-RPC ID",
)
})?;
let params = message.get("params").ok_or_else(|| {
Error::new(ErrorKind::Protocol, "Codex tool request omitted params")
})?;
let call_id = required_string(params, "callId", "Codex tool request")?;
let call = DynamicToolCall {
call_id: call_id.clone(),
tool: required_string(params, "tool", "Codex tool request")?,
arguments: params.get("arguments").cloned().unwrap_or(Value::Null),
};
events
.send(Ok(AgentEvent::ToolCall(call)))
.await
.map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
let (response_call_id, result) = tool_results.recv().await.ok_or_else(|| {
Error::new(ErrorKind::Cancelled, "tool result channel closed")
})?;
if response_call_id != call_id {
return Err(Error::new(
ErrorKind::InvalidInput,
format!(
"tool result call ID '{response_call_id}' does not match pending call '{call_id}'"
),
));
}
write_record(
&mut writer,
events,
json!({
"id":id,
"result":{
"success":result.success,
"contentItems":[{"type":"inputText","text":result.text}],
}
}),
)
.await?;
}
Some("item/completed") => {
if let Some(item) = message.pointer("/params/item")
&& item.get("type").and_then(Value::as_str) == Some("agentMessage")
&& item.get("phase").and_then(Value::as_str) != Some("commentary")
&& let Some(text) = item.get("text").and_then(Value::as_str)
{
answer = text.to_owned();
}
}
Some("thread/tokenUsage/updated") => {
usage = parse_usage(message.pointer("/params/tokenUsage"));
if let Some(usage) = &usage {
events
.send(Ok(AgentEvent::UsageUpdated(usage.clone())))
.await
.map_err(|_| {
Error::new(ErrorKind::Cancelled, "turn event receiver closed")
})?;
}
}
Some("turn/completed") => {
if message.pointer("/params/turn/id").and_then(Value::as_str)
!= Some(turn_id.as_str())
{
continue;
}
let status = message
.pointer("/params/turn/status")
.and_then(Value::as_str)
.unwrap_or("failed");
if status != "completed" {
let detail = message
.pointer("/params/turn/error/message")
.and_then(Value::as_str)
.unwrap_or("Codex turn did not complete");
return Err(Error::new(ErrorKind::Protocol, detail));
}
events
.send(Ok(AgentEvent::Completed(CompletedTurn {
thread_id,
turn_id,
answer,
usage,
})))
.await
.map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
let _ = child.kill().await;
return Ok(());
}
_ => {
if message.get("id").is_some() && message.get("method").is_some() {
let id = message.get("id").cloned().unwrap_or(Value::Null);
write_record(
&mut writer,
events,
json!({"id":id,"error":{"code":-32601,"message":"Method is not supported by this client"}}),
)
.await?;
}
}
}
}
Err(Error::new(
ErrorKind::Unavailable,
"Codex app-server closed before completing the turn",
))
}
fn turn_input(text: &str, images: &[ImageInput]) -> Vec<Value> {
let mut input = Vec::with_capacity(images.len() + 1);
for image in images {
let encoded = BASE64_STANDARD.encode(image.bytes());
input.push(json!({
"type":"image",
"url":format!("data:{};base64,{encoded}", image.media_type().mime_type()),
}));
}
input.push(json!({"type":"text","text":text}));
input
}
async fn write_record<W: AsyncWrite + Unpin>(
writer: &mut W,
events: &mpsc::Sender<Result<AgentEvent>>,
value: Value,
) -> Result<()> {
let mut exact = serde_json::to_string(&value).map_err(|_| {
Error::new(
ErrorKind::InvalidInput,
"Codex request could not be encoded",
)
})?;
exact.push('\n');
writer.write_all(exact.as_bytes()).await.map_err(|_| {
Error::new(
ErrorKind::Unavailable,
"Codex app-server closed its standard input",
)
})?;
writer.flush().await.map_err(|_| {
Error::new(
ErrorKind::Unavailable,
"Codex app-server input could not be flushed",
)
})?;
events
.send(Ok(AgentEvent::ProviderInput(exact)))
.await
.map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))
}
async fn read_response<R: tokio::io::AsyncBufRead + Unpin>(
lines: &mut tokio::io::Lines<R>,
expected_id: u64,
) -> Result<Value> {
while let Some(line) = lines.next_line().await.map_err(|_| {
Error::new(
ErrorKind::Unavailable,
"Codex app-server output could not be read",
)
})? {
let message: Value = serde_json::from_str(&line)
.map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
if message.get("id").and_then(Value::as_u64) == Some(expected_id) {
return Ok(message);
}
}
Err(Error::new(
ErrorKind::Unavailable,
"Codex app-server closed during startup",
))
}
fn require_result<'a>(message: &'a Value, method: &str) -> Result<&'a Value> {
if let Some(error) = message.get("error") {
let detail = error
.get("message")
.and_then(Value::as_str)
.unwrap_or("unknown protocol error");
return Err(Error::new(
ErrorKind::Protocol,
format!("Codex {method} failed: {detail}"),
));
}
message.get("result").ok_or_else(|| {
Error::new(
ErrorKind::Protocol,
format!("Codex {method} response omitted result"),
)
})
}
fn required_string(value: &Value, key: &str, label: &str) -> Result<String> {
value
.get(key)
.and_then(Value::as_str)
.map(str::to_owned)
.ok_or_else(|| Error::new(ErrorKind::Protocol, format!("{label} omitted {key}")))
}
fn parse_usage(value: Option<&Value>) -> Option<TokenUsage> {
let value = value?;
let total = value.get("total")?;
let last = value.get("last");
Some(TokenUsage {
input_tokens: nonnegative(total.get("inputTokens")),
output_tokens: nonnegative(total.get("outputTokens")),
cached_input_tokens: nonnegative(total.get("cachedInputTokens")),
reasoning_output_tokens: nonnegative(total.get("reasoningOutputTokens")),
last_input_tokens: last.map(|value| nonnegative(value.get("inputTokens"))),
last_output_tokens: last.map(|value| nonnegative(value.get("outputTokens"))),
})
}
fn nonnegative(value: Option<&Value>) -> u64 {
value.and_then(Value::as_i64).unwrap_or_default().max(0) as u64
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{DynamicTool, ImageMediaType};
#[test]
fn validates_tool_names_and_schema() {
let mut request = AgentRequest::new("hello", "model");
request.tools = vec![
DynamicTool::new("call", "Call it", json!({"type":"object"})),
DynamicTool::new("call", "Call it again", json!({"type":"object"})),
];
assert_eq!(
validate_request(&request).unwrap_err().kind(),
ErrorKind::InvalidInput
);
}
#[test]
fn validates_image_count() {
assert_eq!(
validate_images(&[]).unwrap_err().kind(),
ErrorKind::InvalidInput
);
let image = ImageInput::new(ImageMediaType::Png, [1]).unwrap();
let images = vec![image; MAX_IMAGE_COUNT + 1];
assert_eq!(
validate_images(&images).unwrap_err().kind(),
ErrorKind::InvalidInput
);
}
#[test]
fn serializes_inline_images_in_order_before_exact_text() {
let images = vec![
ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
ImageInput::new(ImageMediaType::Jpeg, [3, 4]).unwrap(),
];
assert_eq!(
turn_input("Exact prompt\nunchanged", &images),
vec![
json!({"type":"image","url":"data:image/png;base64,AAEC"}),
json!({"type":"image","url":"data:image/jpeg;base64,AwQ="}),
json!({"type":"text","text":"Exact prompt\nunchanged"}),
]
);
assert_eq!(
turn_input("text only", &[]),
vec![json!({"type":"text","text":"text only"})]
);
}
#[test]
fn parses_cumulative_and_last_usage() {
let usage = parse_usage(Some(&json!({
"total":{"inputTokens":20,"outputTokens":7,"cachedInputTokens":4,"reasoningOutputTokens":2},
"last":{"inputTokens":8,"outputTokens":3,"cachedInputTokens":1,"reasoningOutputTokens":1}
})))
.unwrap();
assert_eq!(usage.input_tokens, 20);
assert_eq!(usage.last_output_tokens, Some(3));
}
#[test]
fn app_server_overrides_tool_output_token_limit() {
let command = app_server_command(&CodexConfig::default());
let arguments = command
.as_std()
.get_args()
.map(|argument| argument.to_string_lossy().into_owned())
.collect::<Vec<_>>();
assert!(arguments.windows(2).any(|arguments| {
arguments[0] == "-c"
&& arguments[1] == format!("tool_output_token_limit={TOOL_OUTPUT_TOKEN_LIMIT}")
}));
}
#[cfg(unix)]
#[tokio::test]
async fn completes_a_native_dynamic_tool_turn() {
use std::os::unix::fs::PermissionsExt;
let path = std::env::temp_dir().join(format!(
"kcode-codex-runtime-v2-test-{}",
std::process::id()
));
std::fs::write(
&path,
r##"#!/bin/sh
if [ "$1" = "login" ]; then
echo "Logged in using ChatGPT"
exit 0
fi
read initialize
echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
read initialized
read thread_start
echo '{"id":2,"result":{"thread":{"id":"thread-1"}}}'
read turn_start
echo '{"id":3,"result":{"turn":{"id":"turn-1"}}}'
echo '{"id":77,"method":"item/tool/call","params":{"threadId":"thread-1","turnId":"turn-1","callId":"call-1","tool":"call_ktool","arguments":{"name":"LoadNode","arguments":{"identifier":3}}}}'
read tool_result
echo '{"method":"item/completed","params":{"threadId":"thread-1","turnId":"turn-1","completedAtMs":1,"item":{"id":"message-1","type":"agentMessage","phase":"final_answer","text":"Finished."}}}'
echo '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-1","turnId":"turn-1","tokenUsage":{"total":{"inputTokens":12,"outputTokens":4,"cachedInputTokens":2,"reasoningOutputTokens":1,"totalTokens":16},"last":{"inputTokens":12,"outputTokens":4,"cachedInputTokens":2,"reasoningOutputTokens":1,"totalTokens":16}}}}'
echo '{"method":"turn/completed","params":{"threadId":"thread-1","turn":{"id":"turn-1","items":[],"status":"completed"}}}'
"##,
)
.unwrap();
let mut permissions = std::fs::metadata(&path).unwrap().permissions();
permissions.set_mode(0o700);
std::fs::set_permissions(&path, permissions).unwrap();
let config = CodexConfig {
executable: path.to_string_lossy().into_owned(),
..CodexConfig::default()
};
let codex = Codex::open(config).await.unwrap();
let mut request = AgentRequest::new("Exact input", "test-model");
request.tools.push(DynamicTool::new(
"call_ktool",
"Call one tool",
json!({"type":"object"}),
));
let mut turn = codex.start_turn(request).await.unwrap();
let mut provider_inputs = Vec::new();
let mut usage_updates = Vec::new();
let completed = loop {
match turn.next_event().await.unwrap().unwrap() {
AgentEvent::ProviderInput(exact) => {
provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
}
AgentEvent::UsageUpdated(usage) => usage_updates.push(usage),
AgentEvent::ToolCall(call) => {
assert_eq!(call.arguments["name"], "LoadNode");
turn.respond(call.call_id, ToolResult::success("loaded"))
.await
.unwrap();
}
AgentEvent::Completed(completed) => break completed,
}
};
assert_eq!(completed.answer, "Finished.");
assert_eq!(completed.usage.unwrap().input_tokens, 12);
assert_eq!(usage_updates.len(), 1);
assert_eq!(usage_updates[0].last_input_tokens, Some(12));
let turn_start = provider_inputs
.iter()
.find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
.unwrap();
assert_eq!(
turn_start.pointer("/params/input").unwrap(),
&json!([{"type":"text","text":"Exact input"}])
);
assert!(
provider_inputs
.iter()
.any(|value| value.pointer("/result/success") == Some(&Value::Bool(true)))
);
std::fs::remove_file(path).unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn completes_a_fresh_tool_free_inline_image_turn() {
use std::os::unix::fs::PermissionsExt;
let path = std::env::temp_dir().join(format!(
"kcode-codex-runtime-v2-image-test-{}",
std::process::id()
));
std::fs::write(
&path,
r##"#!/bin/sh
if [ "$1" = "login" ]; then
echo "Logged in using ChatGPT"
exit 0
fi
read initialize
echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
read initialized
read thread_start
echo '{"id":2,"result":{"thread":{"id":"image-thread"}}}'
read turn_start
echo '{"id":3,"result":{"turn":{"id":"image-turn"}}}'
echo '{"method":"item/completed","params":{"threadId":"image-thread","turnId":"image-turn","completedAtMs":1,"item":{"id":"message-1","type":"agentMessage","phase":"final_answer","text":"I see the expected detail."}}}'
echo '{"method":"turn/completed","params":{"threadId":"image-thread","turn":{"id":"image-turn","items":[],"status":"completed"}}}'
"##,
)
.unwrap();
let mut permissions = std::fs::metadata(&path).unwrap().permissions();
permissions.set_mode(0o700);
std::fs::set_permissions(&path, permissions).unwrap();
let config = CodexConfig {
executable: path.to_string_lossy().into_owned(),
..CodexConfig::default()
};
let codex = Codex::open(config).await.unwrap();
let request = ImageTurnRequest::new(
"Read this exactly.",
"test-model",
vec![
ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
ImageInput::new(ImageMediaType::Webp, [3, 4]).unwrap(),
],
);
let mut turn = codex.start_image_turn(request).await.unwrap();
let mut provider_inputs = Vec::new();
let completed = loop {
match turn.next_event().await.unwrap().unwrap() {
AgentEvent::ProviderInput(exact) => {
provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
}
AgentEvent::UsageUpdated(_) => {}
AgentEvent::ToolCall(_) => panic!("image turn exposed a dynamic tool"),
AgentEvent::Completed(completed) => break completed,
}
};
assert_eq!(completed.answer, "I see the expected detail.");
let thread_start = provider_inputs
.iter()
.find(|value| value.get("method").and_then(Value::as_str) == Some("thread/start"))
.unwrap();
assert_eq!(
thread_start.pointer("/params/ephemeral"),
Some(&Value::Bool(true))
);
assert_eq!(
thread_start.pointer("/params/dynamicTools"),
Some(&json!([]))
);
let turn_start = provider_inputs
.iter()
.find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
.unwrap();
assert_eq!(
turn_start.pointer("/params/input").unwrap(),
&json!([
{"type":"image","url":"data:image/png;base64,AAEC"},
{"type":"image","url":"data:image/webp;base64,AwQ="},
{"type":"text","text":"Read this exactly."},
])
);
std::fs::remove_file(path).unwrap();
}
}