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