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, 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(30);
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 let turn_response = read_response(&mut lines, 3).await?;
482 let turn_result = require_result(&turn_response, "turn/start")?;
483 let turn_id = turn_result
484 .pointer("/turn/id")
485 .and_then(Value::as_str)
486 .ok_or_else(|| Error::new(ErrorKind::Protocol, "Codex omitted its turn ID"))?
487 .to_owned();
488
489 let mut answer = String::new();
490 let mut usage = None;
491 while let Some(line) = lines.next_line().await.map_err(|_| {
492 Error::new(
493 ErrorKind::Unavailable,
494 "Codex app-server output could not be read",
495 )
496 })? {
497 let message: Value = serde_json::from_str(&line)
498 .map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
499 match message.get("method").and_then(Value::as_str) {
500 Some("item/tool/call") => {
501 let id = message.get("id").cloned().ok_or_else(|| {
502 Error::new(
503 ErrorKind::Protocol,
504 "Codex tool request omitted its JSON-RPC ID",
505 )
506 })?;
507 let params = message.get("params").ok_or_else(|| {
508 Error::new(ErrorKind::Protocol, "Codex tool request omitted params")
509 })?;
510 let call_id = required_string(params, "callId", "Codex tool request")?;
511 let call = DynamicToolCall {
512 call_id: call_id.clone(),
513 tool: required_string(params, "tool", "Codex tool request")?,
514 arguments: params.get("arguments").cloned().unwrap_or(Value::Null),
515 };
516 events
517 .send(Ok(AgentEvent::ToolCall(call)))
518 .await
519 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
520 let (response_call_id, result) = tool_results.recv().await.ok_or_else(|| {
521 Error::new(ErrorKind::Cancelled, "tool result channel closed")
522 })?;
523 if response_call_id != call_id {
524 return Err(Error::new(
525 ErrorKind::InvalidInput,
526 format!(
527 "tool result call ID '{response_call_id}' does not match pending call '{call_id}'"
528 ),
529 ));
530 }
531 write_record(
532 &mut writer,
533 events,
534 json!({
535 "id":id,
536 "result":{
537 "success":result.success,
538 "contentItems":[{"type":"inputText","text":result.text}],
539 }
540 }),
541 )
542 .await?;
543 }
544 Some("item/completed") => {
545 if let Some(item) = message.pointer("/params/item")
546 && item.get("type").and_then(Value::as_str) == Some("agentMessage")
547 && item.get("phase").and_then(Value::as_str) != Some("commentary")
548 && let Some(text) = item.get("text").and_then(Value::as_str)
549 {
550 answer = text.to_owned();
551 }
552 }
553 Some("thread/tokenUsage/updated") => {
554 usage = parse_usage(message.pointer("/params/tokenUsage"));
555 }
556 Some("turn/completed") => {
557 if message.pointer("/params/turn/id").and_then(Value::as_str)
558 != Some(turn_id.as_str())
559 {
560 continue;
561 }
562 let status = message
563 .pointer("/params/turn/status")
564 .and_then(Value::as_str)
565 .unwrap_or("failed");
566 if status != "completed" {
567 let detail = message
568 .pointer("/params/turn/error/message")
569 .and_then(Value::as_str)
570 .unwrap_or("Codex turn did not complete");
571 return Err(Error::new(ErrorKind::Protocol, detail));
572 }
573 events
574 .send(Ok(AgentEvent::Completed(CompletedTurn {
575 thread_id,
576 turn_id,
577 answer,
578 usage,
579 })))
580 .await
581 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
582 let _ = child.kill().await;
583 return Ok(());
584 }
585 _ => {
586 if message.get("id").is_some() && message.get("method").is_some() {
587 let id = message.get("id").cloned().unwrap_or(Value::Null);
588 write_record(
589 &mut writer,
590 events,
591 json!({"id":id,"error":{"code":-32601,"message":"Method is not supported by this client"}}),
592 )
593 .await?;
594 }
595 }
596 }
597 }
598 Err(Error::new(
599 ErrorKind::Unavailable,
600 "Codex app-server closed before completing the turn",
601 ))
602}
603
604fn turn_input(text: &str, images: &[ImageInput]) -> Vec<Value> {
605 let mut input = Vec::with_capacity(images.len() + 1);
606 for image in images {
607 let encoded = BASE64_STANDARD.encode(image.bytes());
608 input.push(json!({
609 "type":"image",
610 "url":format!("data:{};base64,{encoded}", image.media_type().mime_type()),
611 }));
612 }
613 input.push(json!({"type":"text","text":text}));
614 input
615}
616
617async fn write_record<W: AsyncWrite + Unpin>(
618 writer: &mut W,
619 events: &mpsc::Sender<Result<AgentEvent>>,
620 value: Value,
621) -> Result<()> {
622 let mut exact = serde_json::to_string(&value).map_err(|_| {
623 Error::new(
624 ErrorKind::InvalidInput,
625 "Codex request could not be encoded",
626 )
627 })?;
628 exact.push('\n');
629 writer.write_all(exact.as_bytes()).await.map_err(|_| {
630 Error::new(
631 ErrorKind::Unavailable,
632 "Codex app-server closed its standard input",
633 )
634 })?;
635 writer.flush().await.map_err(|_| {
636 Error::new(
637 ErrorKind::Unavailable,
638 "Codex app-server input could not be flushed",
639 )
640 })?;
641 events
642 .send(Ok(AgentEvent::ProviderInput(exact)))
643 .await
644 .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))
645}
646
647async fn read_response<R: tokio::io::AsyncBufRead + Unpin>(
648 lines: &mut tokio::io::Lines<R>,
649 expected_id: u64,
650) -> Result<Value> {
651 while let Some(line) = lines.next_line().await.map_err(|_| {
652 Error::new(
653 ErrorKind::Unavailable,
654 "Codex app-server output could not be read",
655 )
656 })? {
657 let message: Value = serde_json::from_str(&line)
658 .map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
659 if message.get("id").and_then(Value::as_u64) == Some(expected_id) {
660 return Ok(message);
661 }
662 }
663 Err(Error::new(
664 ErrorKind::Unavailable,
665 "Codex app-server closed during startup",
666 ))
667}
668
669fn require_result<'a>(message: &'a Value, method: &str) -> Result<&'a Value> {
670 if let Some(error) = message.get("error") {
671 let detail = error
672 .get("message")
673 .and_then(Value::as_str)
674 .unwrap_or("unknown protocol error");
675 return Err(Error::new(
676 ErrorKind::Protocol,
677 format!("Codex {method} failed: {detail}"),
678 ));
679 }
680 message.get("result").ok_or_else(|| {
681 Error::new(
682 ErrorKind::Protocol,
683 format!("Codex {method} response omitted result"),
684 )
685 })
686}
687
688fn required_string(value: &Value, key: &str, label: &str) -> Result<String> {
689 value
690 .get(key)
691 .and_then(Value::as_str)
692 .map(str::to_owned)
693 .ok_or_else(|| Error::new(ErrorKind::Protocol, format!("{label} omitted {key}")))
694}
695
696fn parse_usage(value: Option<&Value>) -> Option<TokenUsage> {
697 let value = value?;
698 let total = value.get("total")?;
699 let last = value.get("last");
700 Some(TokenUsage {
701 input_tokens: nonnegative(total.get("inputTokens")),
702 output_tokens: nonnegative(total.get("outputTokens")),
703 cached_input_tokens: nonnegative(total.get("cachedInputTokens")),
704 reasoning_output_tokens: nonnegative(total.get("reasoningOutputTokens")),
705 last_input_tokens: last.map(|value| nonnegative(value.get("inputTokens"))),
706 last_output_tokens: last.map(|value| nonnegative(value.get("outputTokens"))),
707 })
708}
709
710fn nonnegative(value: Option<&Value>) -> u64 {
711 value.and_then(Value::as_i64).unwrap_or_default().max(0) as u64
712}
713
714#[cfg(test)]
715mod tests {
716 use super::*;
717 use crate::{DynamicTool, ImageMediaType};
718
719 #[test]
720 fn validates_tool_names_and_schema() {
721 let mut request = AgentRequest::new("hello", "model");
722 request.tools = vec![
723 DynamicTool::new("call", "Call it", json!({"type":"object"})),
724 DynamicTool::new("call", "Call it again", json!({"type":"object"})),
725 ];
726 assert_eq!(
727 validate_request(&request).unwrap_err().kind(),
728 ErrorKind::InvalidInput
729 );
730 }
731
732 #[test]
733 fn validates_image_count() {
734 assert_eq!(
735 validate_images(&[]).unwrap_err().kind(),
736 ErrorKind::InvalidInput
737 );
738
739 let image = ImageInput::new(ImageMediaType::Png, [1]).unwrap();
740 let images = vec![image; MAX_IMAGE_COUNT + 1];
741 assert_eq!(
742 validate_images(&images).unwrap_err().kind(),
743 ErrorKind::InvalidInput
744 );
745 }
746
747 #[test]
748 fn serializes_inline_images_in_order_before_exact_text() {
749 let images = vec![
750 ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
751 ImageInput::new(ImageMediaType::Jpeg, [3, 4]).unwrap(),
752 ];
753 assert_eq!(
754 turn_input("Exact prompt\nunchanged", &images),
755 vec![
756 json!({"type":"image","url":"data:image/png;base64,AAEC"}),
757 json!({"type":"image","url":"data:image/jpeg;base64,AwQ="}),
758 json!({"type":"text","text":"Exact prompt\nunchanged"}),
759 ]
760 );
761 assert_eq!(
762 turn_input("text only", &[]),
763 vec![json!({"type":"text","text":"text only"})]
764 );
765 }
766
767 #[test]
768 fn parses_cumulative_and_last_usage() {
769 let usage = parse_usage(Some(&json!({
770 "total":{"inputTokens":20,"outputTokens":7,"cachedInputTokens":4,"reasoningOutputTokens":2},
771 "last":{"inputTokens":8,"outputTokens":3,"cachedInputTokens":1,"reasoningOutputTokens":1}
772 })))
773 .unwrap();
774 assert_eq!(usage.input_tokens, 20);
775 assert_eq!(usage.last_output_tokens, Some(3));
776 }
777
778 #[test]
779 fn app_server_overrides_tool_output_token_limit() {
780 let command = app_server_command(&CodexConfig::default());
781 let arguments = command
782 .as_std()
783 .get_args()
784 .map(|argument| argument.to_string_lossy().into_owned())
785 .collect::<Vec<_>>();
786
787 assert!(arguments.windows(2).any(|arguments| {
788 arguments[0] == "-c"
789 && arguments[1] == format!("tool_output_token_limit={TOOL_OUTPUT_TOKEN_LIMIT}")
790 }));
791 }
792
793 #[cfg(unix)]
794 #[tokio::test]
795 async fn completes_a_native_dynamic_tool_turn() {
796 use std::os::unix::fs::PermissionsExt;
797
798 let path = std::env::temp_dir().join(format!(
799 "kcode-codex-runtime-v2-test-{}",
800 std::process::id()
801 ));
802 std::fs::write(
803 &path,
804 r##"#!/bin/sh
805if [ "$1" = "login" ]; then
806 echo "Logged in using ChatGPT"
807 exit 0
808fi
809read initialize
810echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
811read initialized
812read thread_start
813echo '{"id":2,"result":{"thread":{"id":"thread-1"}}}'
814read turn_start
815echo '{"id":3,"result":{"turn":{"id":"turn-1"}}}'
816echo '{"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}}}}'
817read tool_result
818echo '{"method":"item/completed","params":{"threadId":"thread-1","turnId":"turn-1","completedAtMs":1,"item":{"id":"message-1","type":"agentMessage","phase":"final_answer","text":"Finished."}}}'
819echo '{"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}}}}'
820echo '{"method":"turn/completed","params":{"threadId":"thread-1","turn":{"id":"turn-1","items":[],"status":"completed"}}}'
821"##,
822 )
823 .unwrap();
824 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
825 permissions.set_mode(0o700);
826 std::fs::set_permissions(&path, permissions).unwrap();
827
828 let config = CodexConfig {
829 executable: path.to_string_lossy().into_owned(),
830 ..CodexConfig::default()
831 };
832 let codex = Codex::open(config).await.unwrap();
833 let mut request = AgentRequest::new("Exact input", "test-model");
834 request.tools.push(DynamicTool::new(
835 "call_ktool",
836 "Call one tool",
837 json!({"type":"object"}),
838 ));
839 let mut turn = codex.start_turn(request).await.unwrap();
840 let mut provider_inputs = Vec::new();
841 let completed = loop {
842 match turn.next_event().await.unwrap().unwrap() {
843 AgentEvent::ProviderInput(exact) => {
844 provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
845 }
846 AgentEvent::ToolCall(call) => {
847 assert_eq!(call.arguments["name"], "LoadNode");
848 turn.respond(call.call_id, ToolResult::success("loaded"))
849 .await
850 .unwrap();
851 }
852 AgentEvent::Completed(completed) => break completed,
853 }
854 };
855 assert_eq!(completed.answer, "Finished.");
856 assert_eq!(completed.usage.unwrap().input_tokens, 12);
857 let turn_start = provider_inputs
858 .iter()
859 .find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
860 .unwrap();
861 assert_eq!(
862 turn_start.pointer("/params/input").unwrap(),
863 &json!([{"type":"text","text":"Exact input"}])
864 );
865 assert!(
866 provider_inputs
867 .iter()
868 .any(|value| value.pointer("/result/success") == Some(&Value::Bool(true)))
869 );
870 std::fs::remove_file(path).unwrap();
871 }
872
873 #[cfg(unix)]
874 #[tokio::test]
875 async fn completes_a_fresh_tool_free_inline_image_turn() {
876 use std::os::unix::fs::PermissionsExt;
877
878 let path = std::env::temp_dir().join(format!(
879 "kcode-codex-runtime-v2-image-test-{}",
880 std::process::id()
881 ));
882 std::fs::write(
883 &path,
884 r##"#!/bin/sh
885if [ "$1" = "login" ]; then
886 echo "Logged in using ChatGPT"
887 exit 0
888fi
889read initialize
890echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
891read initialized
892read thread_start
893echo '{"id":2,"result":{"thread":{"id":"image-thread"}}}'
894read turn_start
895echo '{"id":3,"result":{"turn":{"id":"image-turn"}}}'
896echo '{"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."}}}'
897echo '{"method":"turn/completed","params":{"threadId":"image-thread","turn":{"id":"image-turn","items":[],"status":"completed"}}}'
898"##,
899 )
900 .unwrap();
901 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
902 permissions.set_mode(0o700);
903 std::fs::set_permissions(&path, permissions).unwrap();
904
905 let config = CodexConfig {
906 executable: path.to_string_lossy().into_owned(),
907 ..CodexConfig::default()
908 };
909 let codex = Codex::open(config).await.unwrap();
910 let request = ImageTurnRequest::new(
911 "Read this exactly.",
912 "test-model",
913 vec![
914 ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
915 ImageInput::new(ImageMediaType::Webp, [3, 4]).unwrap(),
916 ],
917 );
918 let mut turn = codex.start_image_turn(request).await.unwrap();
919 let mut provider_inputs = Vec::new();
920 let completed = loop {
921 match turn.next_event().await.unwrap().unwrap() {
922 AgentEvent::ProviderInput(exact) => {
923 provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
924 }
925 AgentEvent::ToolCall(_) => panic!("image turn exposed a dynamic tool"),
926 AgentEvent::Completed(completed) => break completed,
927 }
928 };
929
930 assert_eq!(completed.answer, "I see the expected detail.");
931 let thread_start = provider_inputs
932 .iter()
933 .find(|value| value.get("method").and_then(Value::as_str) == Some("thread/start"))
934 .unwrap();
935 assert_eq!(
936 thread_start.pointer("/params/ephemeral"),
937 Some(&Value::Bool(true))
938 );
939 assert_eq!(
940 thread_start.pointer("/params/dynamicTools"),
941 Some(&json!([]))
942 );
943 let turn_start = provider_inputs
944 .iter()
945 .find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
946 .unwrap();
947 assert_eq!(
948 turn_start.pointer("/params/input").unwrap(),
949 &json!([
950 {"type":"image","url":"data:image/png;base64,AAEC"},
951 {"type":"image","url":"data:image/webp;base64,AwQ="},
952 {"type":"text","text":"Read this exactly."},
953 ])
954 );
955 std::fs::remove_file(path).unwrap();
956 }
957}