1use std::{collections::HashSet, path::PathBuf, process::Stdio, sync::Arc, time::Duration};
2
3use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
4use serde_json::{Value, json};
5use tokio::{
6 io::{AsyncBufReadExt, AsyncWrite, AsyncWriteExt, BufReader},
7 process::{Child, Command},
8 sync::{mpsc, watch},
9};
10
11use crate::{
12 AgentEvent, AgentRequest, CompletedTurn, DynamicToolCall, Error, ErrorKind, ImageInput,
13 ImageTurnRequest, ModelContext, Result, TokenUsage, ToolResult,
14};
15
16pub const DEFAULT_CODEX_EXECUTABLE: &str = "codex-safe";
18
19const TOOL_OUTPUT_TOKEN_LIMIT: usize = 180_000;
20const MAX_INPUT_CHARACTERS: usize = 1_048_576;
21const MAX_IMAGE_COUNT: usize = 8;
22const MAX_IMAGE_BYTES: usize = 20 * 1024 * 1024;
23const STARTUP_TIMEOUT: Duration = Duration::from_secs(225);
24
25#[derive(Clone, Debug, Eq, PartialEq)]
27pub struct CodexConfig {
28 pub executable: String,
30 pub working_directory: PathBuf,
32 pub base_instruction: String,
34 pub model_catalog: Option<PathBuf>,
36}
37
38impl Default for CodexConfig {
39 fn default() -> Self {
40 Self {
41 executable: DEFAULT_CODEX_EXECUTABLE.into(),
42 working_directory: std::env::temp_dir(),
43 base_instruction: String::new(),
44 model_catalog: None,
45 }
46 }
47}
48
49#[derive(Clone, Debug)]
51pub struct Codex {
52 config: Arc<CodexConfig>,
53}
54
55impl Codex {
56 pub async fn open(config: CodexConfig) -> Result<Self> {
58 validate_config(&config)?;
59 validate_chatgpt_login(&config.executable).await?;
60 Ok(Self {
61 config: Arc::new(config),
62 })
63 }
64
65 pub async fn start_turn(&self, request: AgentRequest) -> Result<AgentTurn> {
67 validate_request(&request)?;
68 self.start_validated_turn(request, Vec::new())
69 }
70
71 pub async fn start_image_turn(&self, request: ImageTurnRequest) -> Result<AgentTurn> {
73 let ImageTurnRequest {
74 prompt,
75 model,
76 images,
77 reasoning_effort,
78 timeout,
79 } = request;
80 validate_images(&images)?;
81
82 let request = AgentRequest {
83 input: prompt,
84 model,
85 reasoning_effort,
86 previous_thread_id: None,
87 tools: Vec::new(),
88 ephemeral: true,
89 timeout,
90 };
91 validate_request(&request)?;
92 self.start_validated_turn(request, images)
93 }
94
95 fn start_validated_turn(
96 &self,
97 request: AgentRequest,
98 images: Vec<ImageInput>,
99 ) -> Result<AgentTurn> {
100 let child = spawn_app_server(&self.config)?;
101 let (event_sender, events) = mpsc::channel(32);
102 let (tool_results, result_receiver) = mpsc::channel(8);
103 let (cancel, cancelled) = watch::channel(false);
104 let config = self.config.clone();
105 let timeout = request.timeout;
106 tokio::spawn(async move {
107 let result = tokio::select! {
108 _ = cancellation(cancelled) => Err(Error::new(ErrorKind::Cancelled, "Codex turn was cancelled")),
109 result = tokio::time::timeout(timeout, run_protocol(child, &config, request, images, &event_sender, result_receiver)) => {
110 match result {
111 Ok(result) => result,
112 Err(_) => Err(Error::new(ErrorKind::Timeout, "Codex turn timed out")),
113 }
114 }
115 };
116 if let Err(error) = result {
117 let _ = event_sender.send(Err(error)).await;
118 }
119 });
120 Ok(AgentTurn {
121 events,
122 tool_results,
123 cancel,
124 })
125 }
126}
127
128pub struct AgentTurn {
130 events: mpsc::Receiver<Result<AgentEvent>>,
131 tool_results: mpsc::Sender<(String, ToolResult)>,
132 cancel: watch::Sender<bool>,
133}
134
135impl AgentTurn {
136 pub async fn next_event(&mut self) -> Option<Result<AgentEvent>> {
138 self.events.recv().await
139 }
140
141 pub async fn respond(&self, call_id: impl Into<String>, result: ToolResult) -> Result<()> {
143 self.tool_results
144 .send((call_id.into(), result))
145 .await
146 .map_err(|_| Error::new(ErrorKind::Cancelled, "Codex turn is no longer running"))
147 }
148
149 pub fn cancel(&self) {
151 let _ = self.cancel.send(true);
152 }
153}
154
155impl Drop for AgentTurn {
156 fn drop(&mut self) {
157 let _ = self.cancel.send(true);
158 }
159}
160
161async fn cancellation(mut cancelled: watch::Receiver<bool>) {
162 if *cancelled.borrow() {
163 return;
164 }
165 while cancelled.changed().await.is_ok() {
166 if *cancelled.borrow() {
167 return;
168 }
169 }
170}
171
172fn validate_config(config: &CodexConfig) -> Result<()> {
173 if config.executable.trim().is_empty() {
174 return Err(Error::new(
175 ErrorKind::InvalidInput,
176 "Codex executable must not be empty",
177 ));
178 }
179 if config.base_instruction.chars().count() > MAX_INPUT_CHARACTERS {
180 return Err(Error::new(
181 ErrorKind::InvalidInput,
182 "Codex base instruction is too large",
183 ));
184 }
185 Ok(())
186}
187
188fn validate_request(request: &AgentRequest) -> Result<()> {
189 if request.input.trim().is_empty() {
190 return Err(Error::new(
191 ErrorKind::InvalidInput,
192 "Codex turn input must not be empty",
193 ));
194 }
195 if request.input.chars().count() > MAX_INPUT_CHARACTERS {
196 return Err(Error::new(
197 ErrorKind::InvalidInput,
198 format!("Codex turn input exceeds {MAX_INPUT_CHARACTERS} characters"),
199 ));
200 }
201 if request.model.trim().is_empty() {
202 return Err(Error::new(
203 ErrorKind::InvalidInput,
204 "Codex model must not be empty",
205 ));
206 }
207 if request.timeout.is_zero() {
208 return Err(Error::new(
209 ErrorKind::InvalidInput,
210 "Codex timeout must be greater than zero",
211 ));
212 }
213 let mut names = HashSet::new();
214 for tool in &request.tools {
215 if tool.name.trim().is_empty() || tool.description.trim().is_empty() {
216 return Err(Error::new(
217 ErrorKind::InvalidInput,
218 "dynamic tool names and descriptions must not be empty",
219 ));
220 }
221 if !tool.input_schema.is_object() {
222 return Err(Error::new(
223 ErrorKind::InvalidInput,
224 format!(
225 "dynamic tool '{}' must use an object JSON Schema",
226 tool.name
227 ),
228 ));
229 }
230 if !names.insert(tool.name.as_str()) {
231 return Err(Error::new(
232 ErrorKind::InvalidInput,
233 format!("dynamic tool '{}' is duplicated", tool.name),
234 ));
235 }
236 }
237 Ok(())
238}
239
240fn validate_images(images: &[ImageInput]) -> Result<()> {
241 if images.is_empty() {
242 return Err(Error::new(
243 ErrorKind::InvalidInput,
244 "Codex image turn requires at least one image",
245 ));
246 }
247 if images.len() > MAX_IMAGE_COUNT {
248 return Err(Error::new(
249 ErrorKind::InvalidInput,
250 format!("Codex image turn supports at most {MAX_IMAGE_COUNT} images"),
251 ));
252 }
253 let total_bytes = images
254 .iter()
255 .try_fold(0usize, |total, image| {
256 total.checked_add(image.bytes().len())
257 })
258 .ok_or_else(|| {
259 Error::new(
260 ErrorKind::InvalidInput,
261 "Codex image input byte count overflowed",
262 )
263 })?;
264 if total_bytes > MAX_IMAGE_BYTES {
265 return Err(Error::new(
266 ErrorKind::InvalidInput,
267 format!("Codex image turn exceeds the {MAX_IMAGE_BYTES}-byte aggregate image limit"),
268 ));
269 }
270 Ok(())
271}
272
273async fn validate_chatgpt_login(executable: &str) -> Result<()> {
274 let output = tokio::time::timeout(
275 STARTUP_TIMEOUT,
276 Command::new(executable)
277 .args(["login", "status"])
278 .env_remove("OPENAI_API_KEY")
279 .env_remove("CODEX_API_KEY")
280 .output(),
281 )
282 .await
283 .map_err(|_| Error::new(ErrorKind::Unavailable, "Codex login check timed out"))?
284 .map_err(|_| {
285 Error::new(
286 ErrorKind::Unavailable,
287 format!("Codex sandbox launcher '{executable}' could not be started"),
288 )
289 })?;
290 let status = format!(
291 "{}\n{}",
292 String::from_utf8_lossy(&output.stdout),
293 String::from_utf8_lossy(&output.stderr)
294 );
295 if !output.status.success() || !status.to_ascii_lowercase().contains("chatgpt") {
296 return Err(Error::new(
297 ErrorKind::Authentication,
298 format!("'{executable}' must be logged in with ChatGPT"),
299 ));
300 }
301 Ok(())
302}
303
304fn spawn_app_server(config: &CodexConfig) -> Result<Child> {
305 app_server_command(config).spawn().map_err(|_| {
306 Error::new(
307 ErrorKind::Unavailable,
308 format!(
309 "Codex sandbox launcher '{}' could not start app-server",
310 config.executable
311 ),
312 )
313 })
314}
315
316fn app_server_command(config: &CodexConfig) -> Command {
317 let mut command = Command::new(&config.executable);
318 command
319 .arg("-c")
320 .arg("web_search=\"disabled\"")
321 .arg("-c")
322 .arg("mcp_servers={}")
323 .arg("-c")
324 .arg("features.shell_tool=false")
325 .arg("-c")
326 .arg("features.apps=false")
327 .arg("-c")
328 .arg("features.browser_use=false")
329 .arg("-c")
330 .arg("features.computer_use=false")
331 .arg("-c")
332 .arg("features.goals=false")
333 .arg("-c")
334 .arg("features.hooks=false")
335 .arg("-c")
336 .arg("features.image_generation=false")
337 .arg("-c")
338 .arg("features.multi_agent=false")
339 .arg("-c")
340 .arg("features.plugins=false")
341 .arg("-c")
342 .arg("features.tool_suggest=false")
343 .arg("-c")
344 .arg("features.remote_plugin=false")
345 .arg("-c")
346 .arg("model_auto_compact_token_limit=9223372036854775807")
347 .arg("-c")
348 .arg(format!("tool_output_token_limit={TOOL_OUTPUT_TOKEN_LIMIT}"));
349 if let Some(path) = &config.model_catalog {
350 command.arg("-c").arg(format!(
351 "model_catalog_json={}",
352 serde_json::to_string(path.to_string_lossy().as_ref())
353 .expect("serializing a path string cannot fail")
354 ));
355 }
356 command
357 .arg("app-server")
358 .arg("--stdio")
359 .current_dir(&config.working_directory)
360 .stdin(Stdio::piped())
361 .stdout(Stdio::piped())
362 .stderr(Stdio::null())
363 .env_remove("OPENAI_API_KEY")
364 .env_remove("CODEX_API_KEY")
365 .kill_on_drop(true);
366 command
367}
368
369async fn run_protocol(
370 mut child: Child,
371 config: &CodexConfig,
372 request: AgentRequest,
373 images: Vec<ImageInput>,
374 events: &mpsc::Sender<Result<AgentEvent>>,
375 mut tool_results: mpsc::Receiver<(String, ToolResult)>,
376) -> Result<()> {
377 let stdin = child.stdin.take().ok_or_else(|| {
378 Error::new(
379 ErrorKind::Unavailable,
380 "Codex app-server standard input is unavailable",
381 )
382 })?;
383 let stdout = child.stdout.take().ok_or_else(|| {
384 Error::new(
385 ErrorKind::Unavailable,
386 "Codex app-server standard output is unavailable",
387 )
388 })?;
389 let mut writer = stdin;
390 let mut lines = BufReader::new(stdout).lines();
391
392 write_record(
393 &mut writer,
394 events,
395 json!({
396 "id":1,
397 "method":"initialize",
398 "params":{
399 "clientInfo":{"name":"kcode-codex-runtime-v2","version":env!("CARGO_PKG_VERSION")},
400 "capabilities":{"experimentalApi":true}
401 }
402 }),
403 )
404 .await?;
405 let initialized = read_response(&mut lines, 1).await?;
406 require_result(&initialized, "initialize")?;
407 write_record(
408 &mut writer,
409 events,
410 json!({"method":"initialized","params":{}}),
411 )
412 .await?;
413
414 let thread_method;
415 let thread_params;
416 if let Some(thread_id) = request.previous_thread_id.as_deref() {
417 thread_method = "thread/resume";
418 thread_params = json!({
419 "threadId":thread_id,
420 "model":request.model,
421 "cwd":config.working_directory,
422 "approvalPolicy":"never",
423 "sandbox":"read-only",
424 "baseInstructions":config.base_instruction,
425 });
426 } else {
427 thread_method = "thread/start";
428 let tools = request
429 .tools
430 .iter()
431 .map(|tool| {
432 json!({
433 "type":"function",
434 "name":tool.name,
435 "description":tool.description,
436 "inputSchema":tool.input_schema,
437 })
438 })
439 .collect::<Vec<_>>();
440 thread_params = json!({
441 "model":request.model,
442 "cwd":config.working_directory,
443 "approvalPolicy":"never",
444 "sandbox":"read-only",
445 "baseInstructions":config.base_instruction,
446 "developerInstructions":"",
447 "dynamicTools":tools,
448 "ephemeral":request.ephemeral,
449 "environments":[],
450 });
451 }
452 write_record(
453 &mut writer,
454 events,
455 json!({"id":2,"method":thread_method,"params":thread_params}),
456 )
457 .await?;
458 let thread_response = read_response(&mut lines, 2).await?;
459 let thread_result = require_result(&thread_response, thread_method)?;
460 let thread_id = thread_result
461 .pointer("/thread/id")
462 .and_then(Value::as_str)
463 .ok_or_else(|| Error::new(ErrorKind::Protocol, "Codex omitted its thread ID"))?
464 .to_owned();
465
466 write_record(
467 &mut writer,
468 events,
469 json!({
470 "id":3,
471 "method":"turn/start",
472 "params":{
473 "threadId":thread_id,
474 "input":turn_input(&request.input, &images),
475 "effort":request.reasoning_effort.as_str(),
476 "approvalPolicy":"never",
477 }
478 }),
479 )
480 .await?;
481 events
482 .send(Ok(AgentEvent::ModelContextSubmitted(ModelContext {
483 input: request.input.clone(),
484 provider: "codex".into(),
485 model: request.model.clone(),
486 reasoning_effort: request.reasoning_effort.as_str().into(),
487 base_instructions: Some(config.base_instruction.clone()),
488 developer_instructions: request.previous_thread_id.is_none().then(String::new),
489 tools: request.tools.clone(),
490 })))
491 .await
492 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
493 let turn_response = read_response(&mut lines, 3).await?;
494 let turn_result = require_result(&turn_response, "turn/start")?;
495 let turn_id = turn_result
496 .pointer("/turn/id")
497 .and_then(Value::as_str)
498 .ok_or_else(|| Error::new(ErrorKind::Protocol, "Codex omitted its turn ID"))?
499 .to_owned();
500
501 let mut answer = String::new();
502 let mut usage = None;
503 while let Some(line) = lines.next_line().await.map_err(|_| {
504 Error::new(
505 ErrorKind::Unavailable,
506 "Codex app-server output could not be read",
507 )
508 })? {
509 let message: Value = serde_json::from_str(&line)
510 .map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
511 match message.get("method").and_then(Value::as_str) {
512 Some("item/tool/call") => {
513 let id = message.get("id").cloned().ok_or_else(|| {
514 Error::new(
515 ErrorKind::Protocol,
516 "Codex tool request omitted its JSON-RPC ID",
517 )
518 })?;
519 let params = message.get("params").ok_or_else(|| {
520 Error::new(ErrorKind::Protocol, "Codex tool request omitted params")
521 })?;
522 let call_id = required_string(params, "callId", "Codex tool request")?;
523 let call = DynamicToolCall {
524 call_id: call_id.clone(),
525 tool: required_string(params, "tool", "Codex tool request")?,
526 arguments: params.get("arguments").cloned().unwrap_or(Value::Null),
527 };
528 events
529 .send(Ok(AgentEvent::ToolCall(call)))
530 .await
531 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
532 let (response_call_id, result) = tool_results.recv().await.ok_or_else(|| {
533 Error::new(ErrorKind::Cancelled, "tool result channel closed")
534 })?;
535 if response_call_id != call_id {
536 return Err(Error::new(
537 ErrorKind::InvalidInput,
538 format!(
539 "tool result call ID '{response_call_id}' does not match pending call '{call_id}'"
540 ),
541 ));
542 }
543 write_record(
544 &mut writer,
545 events,
546 json!({
547 "id":id,
548 "result":{
549 "success":result.success,
550 "contentItems":[{"type":"inputText","text":result.text}],
551 }
552 }),
553 )
554 .await?;
555 }
556 Some("item/completed") => {
557 if let Some(item) = message.pointer("/params/item")
558 && item.get("type").and_then(Value::as_str) == Some("agentMessage")
559 && item.get("phase").and_then(Value::as_str) != Some("commentary")
560 && let Some(text) = item.get("text").and_then(Value::as_str)
561 {
562 answer = text.to_owned();
563 }
564 }
565 Some("thread/tokenUsage/updated") => {
566 usage = parse_usage(message.pointer("/params/tokenUsage"));
567 if let Some(usage) = &usage {
568 events
569 .send(Ok(AgentEvent::UsageUpdated(usage.clone())))
570 .await
571 .map_err(|_| {
572 Error::new(ErrorKind::Cancelled, "turn event receiver closed")
573 })?;
574 }
575 }
576 Some("turn/completed") => {
577 if message.pointer("/params/turn/id").and_then(Value::as_str)
578 != Some(turn_id.as_str())
579 {
580 continue;
581 }
582 let status = message
583 .pointer("/params/turn/status")
584 .and_then(Value::as_str)
585 .unwrap_or("failed");
586 if status != "completed" {
587 let detail = message
588 .pointer("/params/turn/error/message")
589 .and_then(Value::as_str)
590 .unwrap_or("Codex turn did not complete");
591 return Err(Error::new(ErrorKind::Protocol, detail));
592 }
593 events
594 .send(Ok(AgentEvent::Completed(CompletedTurn {
595 thread_id,
596 turn_id,
597 answer,
598 usage,
599 })))
600 .await
601 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
602 let _ = child.kill().await;
603 return Ok(());
604 }
605 _ => {
606 if message.get("id").is_some() && message.get("method").is_some() {
607 let id = message.get("id").cloned().unwrap_or(Value::Null);
608 write_record(
609 &mut writer,
610 events,
611 json!({"id":id,"error":{"code":-32601,"message":"Method is not supported by this client"}}),
612 )
613 .await?;
614 }
615 }
616 }
617 }
618 Err(Error::new(
619 ErrorKind::Unavailable,
620 "Codex app-server closed before completing the turn",
621 ))
622}
623
624fn turn_input(text: &str, images: &[ImageInput]) -> Vec<Value> {
625 let mut input = Vec::with_capacity(images.len() + 1);
626 for image in images {
627 let encoded = BASE64_STANDARD.encode(image.bytes());
628 input.push(json!({
629 "type":"image",
630 "url":format!("data:{};base64,{encoded}", image.media_type().mime_type()),
631 }));
632 }
633 input.push(json!({"type":"text","text":text}));
634 input
635}
636
637async fn write_record<W: AsyncWrite + Unpin>(
638 writer: &mut W,
639 events: &mpsc::Sender<Result<AgentEvent>>,
640 value: Value,
641) -> Result<()> {
642 let mut exact = serde_json::to_string(&value).map_err(|_| {
643 Error::new(
644 ErrorKind::InvalidInput,
645 "Codex request could not be encoded",
646 )
647 })?;
648 exact.push('\n');
649 writer.write_all(exact.as_bytes()).await.map_err(|_| {
650 Error::new(
651 ErrorKind::Unavailable,
652 "Codex app-server closed its standard input",
653 )
654 })?;
655 writer.flush().await.map_err(|_| {
656 Error::new(
657 ErrorKind::Unavailable,
658 "Codex app-server input could not be flushed",
659 )
660 })?;
661 events
662 .send(Ok(AgentEvent::ProviderInput(exact)))
663 .await
664 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))
665}
666
667async fn read_response<R: tokio::io::AsyncBufRead + Unpin>(
668 lines: &mut tokio::io::Lines<R>,
669 expected_id: u64,
670) -> Result<Value> {
671 while let Some(line) = lines.next_line().await.map_err(|_| {
672 Error::new(
673 ErrorKind::Unavailable,
674 "Codex app-server output could not be read",
675 )
676 })? {
677 let message: Value = serde_json::from_str(&line)
678 .map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
679 if message.get("id").and_then(Value::as_u64) == Some(expected_id) {
680 return Ok(message);
681 }
682 }
683 Err(Error::new(
684 ErrorKind::Unavailable,
685 "Codex app-server closed during startup",
686 ))
687}
688
689fn require_result<'a>(message: &'a Value, method: &str) -> Result<&'a Value> {
690 if let Some(error) = message.get("error") {
691 let detail = error
692 .get("message")
693 .and_then(Value::as_str)
694 .unwrap_or("unknown protocol error");
695 return Err(Error::new(
696 ErrorKind::Protocol,
697 format!("Codex {method} failed: {detail}"),
698 ));
699 }
700 message.get("result").ok_or_else(|| {
701 Error::new(
702 ErrorKind::Protocol,
703 format!("Codex {method} response omitted result"),
704 )
705 })
706}
707
708fn required_string(value: &Value, key: &str, label: &str) -> Result<String> {
709 value
710 .get(key)
711 .and_then(Value::as_str)
712 .map(str::to_owned)
713 .ok_or_else(|| Error::new(ErrorKind::Protocol, format!("{label} omitted {key}")))
714}
715
716fn parse_usage(value: Option<&Value>) -> Option<TokenUsage> {
717 let value = value?;
718 let total = value.get("total")?;
719 let last = value.get("last");
720 Some(TokenUsage {
721 input_tokens: nonnegative(total.get("inputTokens")),
722 output_tokens: nonnegative(total.get("outputTokens")),
723 cached_input_tokens: nonnegative(total.get("cachedInputTokens")),
724 reasoning_output_tokens: nonnegative(total.get("reasoningOutputTokens")),
725 last_input_tokens: last.map(|value| nonnegative(value.get("inputTokens"))),
726 last_output_tokens: last.map(|value| nonnegative(value.get("outputTokens"))),
727 })
728}
729
730fn nonnegative(value: Option<&Value>) -> u64 {
731 value.and_then(Value::as_i64).unwrap_or_default().max(0) as u64
732}
733
734#[cfg(test)]
735mod tests {
736 use super::*;
737 use crate::{DynamicTool, ImageMediaType};
738
739 #[cfg(unix)]
740 static PROCESS_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
741
742 #[test]
743 fn validates_tool_names_and_schema() {
744 let mut request = AgentRequest::new("hello", "model");
745 request.tools = vec![
746 DynamicTool::new("call", "Call it", json!({"type":"object"})),
747 DynamicTool::new("call", "Call it again", json!({"type":"object"})),
748 ];
749 assert_eq!(
750 validate_request(&request).unwrap_err().kind(),
751 ErrorKind::InvalidInput
752 );
753 }
754
755 #[test]
756 fn validates_image_count() {
757 assert_eq!(
758 validate_images(&[]).unwrap_err().kind(),
759 ErrorKind::InvalidInput
760 );
761
762 let image = ImageInput::new(ImageMediaType::Png, [1]).unwrap();
763 let images = vec![image; MAX_IMAGE_COUNT + 1];
764 assert_eq!(
765 validate_images(&images).unwrap_err().kind(),
766 ErrorKind::InvalidInput
767 );
768 }
769
770 #[test]
771 fn serializes_inline_images_in_order_before_exact_text() {
772 let images = vec![
773 ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
774 ImageInput::new(ImageMediaType::Jpeg, [3, 4]).unwrap(),
775 ];
776 assert_eq!(
777 turn_input("Exact prompt\nunchanged", &images),
778 vec![
779 json!({"type":"image","url":"data:image/png;base64,AAEC"}),
780 json!({"type":"image","url":"data:image/jpeg;base64,AwQ="}),
781 json!({"type":"text","text":"Exact prompt\nunchanged"}),
782 ]
783 );
784 assert_eq!(
785 turn_input("text only", &[]),
786 vec![json!({"type":"text","text":"text only"})]
787 );
788 }
789
790 #[test]
791 fn parses_cumulative_and_last_usage() {
792 let usage = parse_usage(Some(&json!({
793 "total":{"inputTokens":20,"outputTokens":7,"cachedInputTokens":4,"reasoningOutputTokens":2},
794 "last":{"inputTokens":8,"outputTokens":3,"cachedInputTokens":1,"reasoningOutputTokens":1}
795 })))
796 .unwrap();
797 assert_eq!(usage.input_tokens, 20);
798 assert_eq!(usage.last_output_tokens, Some(3));
799 }
800
801 #[test]
802 fn app_server_overrides_tool_output_token_limit() {
803 let command = app_server_command(&CodexConfig::default());
804 let arguments = command
805 .as_std()
806 .get_args()
807 .map(|argument| argument.to_string_lossy().into_owned())
808 .collect::<Vec<_>>();
809
810 assert!(arguments.windows(2).any(|arguments| {
811 arguments[0] == "-c"
812 && arguments[1] == format!("tool_output_token_limit={TOOL_OUTPUT_TOKEN_LIMIT}")
813 }));
814 }
815
816 #[cfg(unix)]
817 #[tokio::test]
818 async fn completes_a_native_dynamic_tool_turn() {
819 use std::os::unix::fs::PermissionsExt;
820
821 let _process_test_guard = PROCESS_TEST_LOCK.lock().await;
822
823 let path = std::env::temp_dir().join(format!(
824 "kcode-codex-runtime-v2-test-{}",
825 std::process::id()
826 ));
827 std::fs::write(
828 &path,
829 r##"#!/bin/sh
830if [ "$1" = "login" ]; then
831 echo "Logged in using ChatGPT"
832 exit 0
833fi
834read initialize
835echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
836read initialized
837read thread_start
838echo '{"id":2,"result":{"thread":{"id":"thread-1"}}}'
839read turn_start
840echo '{"id":3,"result":{"turn":{"id":"turn-1"}}}'
841echo '{"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}}}}'
842read tool_result
843echo '{"method":"item/completed","params":{"threadId":"thread-1","turnId":"turn-1","completedAtMs":1,"item":{"id":"message-1","type":"agentMessage","phase":"final_answer","text":"Finished."}}}'
844echo '{"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}}}}'
845echo '{"method":"turn/completed","params":{"threadId":"thread-1","turn":{"id":"turn-1","items":[],"status":"completed"}}}'
846"##,
847 )
848 .unwrap();
849 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
850 permissions.set_mode(0o700);
851 std::fs::set_permissions(&path, permissions).unwrap();
852
853 let config = CodexConfig {
854 executable: path.to_string_lossy().into_owned(),
855 ..CodexConfig::default()
856 };
857 let codex = Codex::open(config).await.unwrap();
858 let mut request = AgentRequest::new("Exact input", "test-model");
859 request.tools.push(DynamicTool::new(
860 "call_ktool",
861 "Call one tool",
862 json!({"type":"object"}),
863 ));
864 let mut turn = codex.start_turn(request).await.unwrap();
865 let mut provider_inputs = Vec::new();
866 let mut model_context = None;
867 let mut usage_updates = Vec::new();
868 let completed = loop {
869 match turn.next_event().await.unwrap().unwrap() {
870 AgentEvent::ProviderInput(exact) => {
871 provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
872 }
873 AgentEvent::ModelContextSubmitted(context) => model_context = Some(context),
874 AgentEvent::UsageUpdated(usage) => usage_updates.push(usage),
875 AgentEvent::ToolCall(call) => {
876 assert_eq!(call.arguments["name"], "LoadNode");
877 turn.respond(call.call_id, ToolResult::success("loaded"))
878 .await
879 .unwrap();
880 }
881 AgentEvent::Completed(completed) => break completed,
882 }
883 };
884 assert_eq!(completed.answer, "Finished.");
885 assert_eq!(completed.usage.unwrap().input_tokens, 12);
886 assert_eq!(usage_updates.len(), 1);
887 assert_eq!(usage_updates[0].last_input_tokens, Some(12));
888 let context = model_context.unwrap();
889 assert_eq!(context.input, "Exact input");
890 assert_eq!(context.provider, "codex");
891 assert_eq!(context.model, "test-model");
892 assert_eq!(context.reasoning_effort, "xhigh");
893 assert_eq!(context.base_instructions.as_deref(), Some(""));
894 assert_eq!(context.developer_instructions.as_deref(), Some(""));
895 assert_eq!(context.tools.len(), 1);
896 let turn_start = provider_inputs
897 .iter()
898 .find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
899 .unwrap();
900 assert_eq!(
901 turn_start.pointer("/params/input").unwrap(),
902 &json!([{"type":"text","text":"Exact input"}])
903 );
904 assert!(
905 provider_inputs
906 .iter()
907 .any(|value| value.pointer("/result/success") == Some(&Value::Bool(true)))
908 );
909 std::fs::remove_file(path).unwrap();
910 }
911
912 #[cfg(unix)]
913 #[tokio::test]
914 async fn completes_a_fresh_tool_free_inline_image_turn() {
915 use std::os::unix::fs::PermissionsExt;
916
917 let _process_test_guard = PROCESS_TEST_LOCK.lock().await;
918
919 let path = std::env::temp_dir().join(format!(
920 "kcode-codex-runtime-v2-image-test-{}",
921 std::process::id()
922 ));
923 std::fs::write(
924 &path,
925 r##"#!/bin/sh
926if [ "$1" = "login" ]; then
927 echo "Logged in using ChatGPT"
928 exit 0
929fi
930read initialize
931echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
932read initialized
933read thread_start
934echo '{"id":2,"result":{"thread":{"id":"image-thread"}}}'
935read turn_start
936echo '{"id":3,"result":{"turn":{"id":"image-turn"}}}'
937echo '{"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."}}}'
938echo '{"method":"turn/completed","params":{"threadId":"image-thread","turn":{"id":"image-turn","items":[],"status":"completed"}}}'
939"##,
940 )
941 .unwrap();
942 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
943 permissions.set_mode(0o700);
944 std::fs::set_permissions(&path, permissions).unwrap();
945
946 let config = CodexConfig {
947 executable: path.to_string_lossy().into_owned(),
948 ..CodexConfig::default()
949 };
950 let codex = Codex::open(config).await.unwrap();
951 let request = ImageTurnRequest::new(
952 "Read this exactly.",
953 "test-model",
954 vec![
955 ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
956 ImageInput::new(ImageMediaType::Webp, [3, 4]).unwrap(),
957 ],
958 );
959 let mut turn = codex.start_image_turn(request).await.unwrap();
960 let mut provider_inputs = Vec::new();
961 let mut model_context = None;
962 let completed = loop {
963 match turn.next_event().await.unwrap().unwrap() {
964 AgentEvent::ProviderInput(exact) => {
965 provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
966 }
967 AgentEvent::ModelContextSubmitted(context) => model_context = Some(context),
968 AgentEvent::UsageUpdated(_) => {}
969 AgentEvent::ToolCall(_) => panic!("image turn exposed a dynamic tool"),
970 AgentEvent::Completed(completed) => break completed,
971 }
972 };
973
974 assert_eq!(completed.answer, "I see the expected detail.");
975 let context = model_context.unwrap();
976 assert_eq!(context.input, "Read this exactly.");
977 assert!(context.tools.is_empty());
978 let thread_start = provider_inputs
979 .iter()
980 .find(|value| value.get("method").and_then(Value::as_str) == Some("thread/start"))
981 .unwrap();
982 assert_eq!(
983 thread_start.pointer("/params/ephemeral"),
984 Some(&Value::Bool(true))
985 );
986 assert_eq!(
987 thread_start.pointer("/params/dynamicTools"),
988 Some(&json!([]))
989 );
990 let turn_start = provider_inputs
991 .iter()
992 .find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
993 .unwrap();
994 assert_eq!(
995 turn_start.pointer("/params/input").unwrap(),
996 &json!([
997 {"type":"image","url":"data:image/png;base64,AAEC"},
998 {"type":"image","url":"data:image/webp;base64,AwQ="},
999 {"type":"text","text":"Read this exactly."},
1000 ])
1001 );
1002 std::fs::remove_file(path).unwrap();
1003 }
1004}