Skip to main content

kcode_codex_runtime_v2/
runtime.rs

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
16/// Default sandboxed Codex launcher.
17pub 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(90);
24
25/// Runtime configuration applied to every Codex app-server process.
26#[derive(Clone, Debug, Eq, PartialEq)]
27pub struct CodexConfig {
28    /// Codex or sandbox-launcher executable.
29    pub executable: String,
30    /// Directory supplied to Codex as its working directory.
31    pub working_directory: PathBuf,
32    /// Fixed base instruction for every fresh or resumed thread.
33    pub base_instruction: String,
34    /// Optional sanitized Codex model-catalog JSON path.
35    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/// Cloneable Codex app-server runtime.
50#[derive(Clone, Debug)]
51pub struct Codex {
52    config: Arc<CodexConfig>,
53}
54
55impl Codex {
56    /// Validates configuration and requires a ChatGPT-authenticated Codex CLI.
57    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    /// Starts one standard Codex turn and returns its interactive event stream.
66    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    /// Starts one fresh, ephemeral, tool-free turn with inline image bytes.
72    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
128/// Interactive handle for a running Codex turn.
129pub 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    /// Waits for the next provider-input, tool-call, or completion event.
137    pub async fn next_event(&mut self) -> Option<Result<AgentEvent>> {
138        self.events.recv().await
139    }
140
141    /// Returns the result for the currently pending dynamic tool call.
142    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    /// Stops this turn. Dropping the handle has the same effect.
150    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                if let Some(usage) = &usage {
556                    events
557                        .send(Ok(AgentEvent::UsageUpdated(usage.clone())))
558                        .await
559                        .map_err(|_| {
560                            Error::new(ErrorKind::Cancelled, "turn event receiver closed")
561                        })?;
562                }
563            }
564            Some("turn/completed") => {
565                if message.pointer("/params/turn/id").and_then(Value::as_str)
566                    != Some(turn_id.as_str())
567                {
568                    continue;
569                }
570                let status = message
571                    .pointer("/params/turn/status")
572                    .and_then(Value::as_str)
573                    .unwrap_or("failed");
574                if status != "completed" {
575                    let detail = message
576                        .pointer("/params/turn/error/message")
577                        .and_then(Value::as_str)
578                        .unwrap_or("Codex turn did not complete");
579                    return Err(Error::new(ErrorKind::Protocol, detail));
580                }
581                events
582                    .send(Ok(AgentEvent::Completed(CompletedTurn {
583                        thread_id,
584                        turn_id,
585                        answer,
586                        usage,
587                    })))
588                    .await
589                    .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))?;
590                let _ = child.kill().await;
591                return Ok(());
592            }
593            _ => {
594                if message.get("id").is_some() && message.get("method").is_some() {
595                    let id = message.get("id").cloned().unwrap_or(Value::Null);
596                    write_record(
597                        &mut writer,
598                        events,
599                        json!({"id":id,"error":{"code":-32601,"message":"Method is not supported by this client"}}),
600                    )
601                    .await?;
602                }
603            }
604        }
605    }
606    Err(Error::new(
607        ErrorKind::Unavailable,
608        "Codex app-server closed before completing the turn",
609    ))
610}
611
612fn turn_input(text: &str, images: &[ImageInput]) -> Vec<Value> {
613    let mut input = Vec::with_capacity(images.len() + 1);
614    for image in images {
615        let encoded = BASE64_STANDARD.encode(image.bytes());
616        input.push(json!({
617            "type":"image",
618            "url":format!("data:{};base64,{encoded}", image.media_type().mime_type()),
619        }));
620    }
621    input.push(json!({"type":"text","text":text}));
622    input
623}
624
625async fn write_record<W: AsyncWrite + Unpin>(
626    writer: &mut W,
627    events: &mpsc::Sender<Result<AgentEvent>>,
628    value: Value,
629) -> Result<()> {
630    let mut exact = serde_json::to_string(&value).map_err(|_| {
631        Error::new(
632            ErrorKind::InvalidInput,
633            "Codex request could not be encoded",
634        )
635    })?;
636    exact.push('\n');
637    writer.write_all(exact.as_bytes()).await.map_err(|_| {
638        Error::new(
639            ErrorKind::Unavailable,
640            "Codex app-server closed its standard input",
641        )
642    })?;
643    writer.flush().await.map_err(|_| {
644        Error::new(
645            ErrorKind::Unavailable,
646            "Codex app-server input could not be flushed",
647        )
648    })?;
649    events
650        .send(Ok(AgentEvent::ProviderInput(exact)))
651        .await
652        .map_err(|_| Error::new(ErrorKind::Cancelled, "turn event receiver closed"))
653}
654
655async fn read_response<R: tokio::io::AsyncBufRead + Unpin>(
656    lines: &mut tokio::io::Lines<R>,
657    expected_id: u64,
658) -> Result<Value> {
659    while let Some(line) = lines.next_line().await.map_err(|_| {
660        Error::new(
661            ErrorKind::Unavailable,
662            "Codex app-server output could not be read",
663        )
664    })? {
665        let message: Value = serde_json::from_str(&line)
666            .map_err(|_| Error::new(ErrorKind::Protocol, "Codex emitted invalid JSONL"))?;
667        if message.get("id").and_then(Value::as_u64) == Some(expected_id) {
668            return Ok(message);
669        }
670    }
671    Err(Error::new(
672        ErrorKind::Unavailable,
673        "Codex app-server closed during startup",
674    ))
675}
676
677fn require_result<'a>(message: &'a Value, method: &str) -> Result<&'a Value> {
678    if let Some(error) = message.get("error") {
679        let detail = error
680            .get("message")
681            .and_then(Value::as_str)
682            .unwrap_or("unknown protocol error");
683        return Err(Error::new(
684            ErrorKind::Protocol,
685            format!("Codex {method} failed: {detail}"),
686        ));
687    }
688    message.get("result").ok_or_else(|| {
689        Error::new(
690            ErrorKind::Protocol,
691            format!("Codex {method} response omitted result"),
692        )
693    })
694}
695
696fn required_string(value: &Value, key: &str, label: &str) -> Result<String> {
697    value
698        .get(key)
699        .and_then(Value::as_str)
700        .map(str::to_owned)
701        .ok_or_else(|| Error::new(ErrorKind::Protocol, format!("{label} omitted {key}")))
702}
703
704fn parse_usage(value: Option<&Value>) -> Option<TokenUsage> {
705    let value = value?;
706    let total = value.get("total")?;
707    let last = value.get("last");
708    Some(TokenUsage {
709        input_tokens: nonnegative(total.get("inputTokens")),
710        output_tokens: nonnegative(total.get("outputTokens")),
711        cached_input_tokens: nonnegative(total.get("cachedInputTokens")),
712        reasoning_output_tokens: nonnegative(total.get("reasoningOutputTokens")),
713        last_input_tokens: last.map(|value| nonnegative(value.get("inputTokens"))),
714        last_output_tokens: last.map(|value| nonnegative(value.get("outputTokens"))),
715    })
716}
717
718fn nonnegative(value: Option<&Value>) -> u64 {
719    value.and_then(Value::as_i64).unwrap_or_default().max(0) as u64
720}
721
722#[cfg(test)]
723mod tests {
724    use super::*;
725    use crate::{DynamicTool, ImageMediaType};
726
727    #[test]
728    fn validates_tool_names_and_schema() {
729        let mut request = AgentRequest::new("hello", "model");
730        request.tools = vec![
731            DynamicTool::new("call", "Call it", json!({"type":"object"})),
732            DynamicTool::new("call", "Call it again", json!({"type":"object"})),
733        ];
734        assert_eq!(
735            validate_request(&request).unwrap_err().kind(),
736            ErrorKind::InvalidInput
737        );
738    }
739
740    #[test]
741    fn validates_image_count() {
742        assert_eq!(
743            validate_images(&[]).unwrap_err().kind(),
744            ErrorKind::InvalidInput
745        );
746
747        let image = ImageInput::new(ImageMediaType::Png, [1]).unwrap();
748        let images = vec![image; MAX_IMAGE_COUNT + 1];
749        assert_eq!(
750            validate_images(&images).unwrap_err().kind(),
751            ErrorKind::InvalidInput
752        );
753    }
754
755    #[test]
756    fn serializes_inline_images_in_order_before_exact_text() {
757        let images = vec![
758            ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
759            ImageInput::new(ImageMediaType::Jpeg, [3, 4]).unwrap(),
760        ];
761        assert_eq!(
762            turn_input("Exact prompt\nunchanged", &images),
763            vec![
764                json!({"type":"image","url":"data:image/png;base64,AAEC"}),
765                json!({"type":"image","url":"data:image/jpeg;base64,AwQ="}),
766                json!({"type":"text","text":"Exact prompt\nunchanged"}),
767            ]
768        );
769        assert_eq!(
770            turn_input("text only", &[]),
771            vec![json!({"type":"text","text":"text only"})]
772        );
773    }
774
775    #[test]
776    fn parses_cumulative_and_last_usage() {
777        let usage = parse_usage(Some(&json!({
778            "total":{"inputTokens":20,"outputTokens":7,"cachedInputTokens":4,"reasoningOutputTokens":2},
779            "last":{"inputTokens":8,"outputTokens":3,"cachedInputTokens":1,"reasoningOutputTokens":1}
780        })))
781        .unwrap();
782        assert_eq!(usage.input_tokens, 20);
783        assert_eq!(usage.last_output_tokens, Some(3));
784    }
785
786    #[test]
787    fn app_server_overrides_tool_output_token_limit() {
788        let command = app_server_command(&CodexConfig::default());
789        let arguments = command
790            .as_std()
791            .get_args()
792            .map(|argument| argument.to_string_lossy().into_owned())
793            .collect::<Vec<_>>();
794
795        assert!(arguments.windows(2).any(|arguments| {
796            arguments[0] == "-c"
797                && arguments[1] == format!("tool_output_token_limit={TOOL_OUTPUT_TOKEN_LIMIT}")
798        }));
799    }
800
801    #[cfg(unix)]
802    #[tokio::test]
803    async fn completes_a_native_dynamic_tool_turn() {
804        use std::os::unix::fs::PermissionsExt;
805
806        let path = std::env::temp_dir().join(format!(
807            "kcode-codex-runtime-v2-test-{}",
808            std::process::id()
809        ));
810        std::fs::write(
811            &path,
812            r##"#!/bin/sh
813if [ "$1" = "login" ]; then
814  echo "Logged in using ChatGPT"
815  exit 0
816fi
817read initialize
818echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
819read initialized
820read thread_start
821echo '{"id":2,"result":{"thread":{"id":"thread-1"}}}'
822read turn_start
823echo '{"id":3,"result":{"turn":{"id":"turn-1"}}}'
824echo '{"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}}}}'
825read tool_result
826echo '{"method":"item/completed","params":{"threadId":"thread-1","turnId":"turn-1","completedAtMs":1,"item":{"id":"message-1","type":"agentMessage","phase":"final_answer","text":"Finished."}}}'
827echo '{"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}}}}'
828echo '{"method":"turn/completed","params":{"threadId":"thread-1","turn":{"id":"turn-1","items":[],"status":"completed"}}}'
829"##,
830        )
831        .unwrap();
832        let mut permissions = std::fs::metadata(&path).unwrap().permissions();
833        permissions.set_mode(0o700);
834        std::fs::set_permissions(&path, permissions).unwrap();
835
836        let config = CodexConfig {
837            executable: path.to_string_lossy().into_owned(),
838            ..CodexConfig::default()
839        };
840        let codex = Codex::open(config).await.unwrap();
841        let mut request = AgentRequest::new("Exact input", "test-model");
842        request.tools.push(DynamicTool::new(
843            "call_ktool",
844            "Call one tool",
845            json!({"type":"object"}),
846        ));
847        let mut turn = codex.start_turn(request).await.unwrap();
848        let mut provider_inputs = Vec::new();
849        let mut usage_updates = Vec::new();
850        let completed = loop {
851            match turn.next_event().await.unwrap().unwrap() {
852                AgentEvent::ProviderInput(exact) => {
853                    provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
854                }
855                AgentEvent::UsageUpdated(usage) => usage_updates.push(usage),
856                AgentEvent::ToolCall(call) => {
857                    assert_eq!(call.arguments["name"], "LoadNode");
858                    turn.respond(call.call_id, ToolResult::success("loaded"))
859                        .await
860                        .unwrap();
861                }
862                AgentEvent::Completed(completed) => break completed,
863            }
864        };
865        assert_eq!(completed.answer, "Finished.");
866        assert_eq!(completed.usage.unwrap().input_tokens, 12);
867        assert_eq!(usage_updates.len(), 1);
868        assert_eq!(usage_updates[0].last_input_tokens, Some(12));
869        let turn_start = provider_inputs
870            .iter()
871            .find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
872            .unwrap();
873        assert_eq!(
874            turn_start.pointer("/params/input").unwrap(),
875            &json!([{"type":"text","text":"Exact input"}])
876        );
877        assert!(
878            provider_inputs
879                .iter()
880                .any(|value| value.pointer("/result/success") == Some(&Value::Bool(true)))
881        );
882        std::fs::remove_file(path).unwrap();
883    }
884
885    #[cfg(unix)]
886    #[tokio::test]
887    async fn completes_a_fresh_tool_free_inline_image_turn() {
888        use std::os::unix::fs::PermissionsExt;
889
890        let path = std::env::temp_dir().join(format!(
891            "kcode-codex-runtime-v2-image-test-{}",
892            std::process::id()
893        ));
894        std::fs::write(
895            &path,
896            r##"#!/bin/sh
897if [ "$1" = "login" ]; then
898  echo "Logged in using ChatGPT"
899  exit 0
900fi
901read initialize
902echo '{"id":1,"result":{"userAgent":"test","platformFamily":"unix","platformOs":"linux","codexHome":"/tmp"}}'
903read initialized
904read thread_start
905echo '{"id":2,"result":{"thread":{"id":"image-thread"}}}'
906read turn_start
907echo '{"id":3,"result":{"turn":{"id":"image-turn"}}}'
908echo '{"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."}}}'
909echo '{"method":"turn/completed","params":{"threadId":"image-thread","turn":{"id":"image-turn","items":[],"status":"completed"}}}'
910"##,
911        )
912        .unwrap();
913        let mut permissions = std::fs::metadata(&path).unwrap().permissions();
914        permissions.set_mode(0o700);
915        std::fs::set_permissions(&path, permissions).unwrap();
916
917        let config = CodexConfig {
918            executable: path.to_string_lossy().into_owned(),
919            ..CodexConfig::default()
920        };
921        let codex = Codex::open(config).await.unwrap();
922        let request = ImageTurnRequest::new(
923            "Read this exactly.",
924            "test-model",
925            vec![
926                ImageInput::new(ImageMediaType::Png, [0, 1, 2]).unwrap(),
927                ImageInput::new(ImageMediaType::Webp, [3, 4]).unwrap(),
928            ],
929        );
930        let mut turn = codex.start_image_turn(request).await.unwrap();
931        let mut provider_inputs = Vec::new();
932        let completed = loop {
933            match turn.next_event().await.unwrap().unwrap() {
934                AgentEvent::ProviderInput(exact) => {
935                    provider_inputs.push(serde_json::from_str::<Value>(&exact).unwrap());
936                }
937                AgentEvent::UsageUpdated(_) => {}
938                AgentEvent::ToolCall(_) => panic!("image turn exposed a dynamic tool"),
939                AgentEvent::Completed(completed) => break completed,
940            }
941        };
942
943        assert_eq!(completed.answer, "I see the expected detail.");
944        let thread_start = provider_inputs
945            .iter()
946            .find(|value| value.get("method").and_then(Value::as_str) == Some("thread/start"))
947            .unwrap();
948        assert_eq!(
949            thread_start.pointer("/params/ephemeral"),
950            Some(&Value::Bool(true))
951        );
952        assert_eq!(
953            thread_start.pointer("/params/dynamicTools"),
954            Some(&json!([]))
955        );
956        let turn_start = provider_inputs
957            .iter()
958            .find(|value| value.get("method").and_then(Value::as_str) == Some("turn/start"))
959            .unwrap();
960        assert_eq!(
961            turn_start.pointer("/params/input").unwrap(),
962            &json!([
963                {"type":"image","url":"data:image/png;base64,AAEC"},
964                {"type":"image","url":"data:image/webp;base64,AwQ="},
965                {"type":"text","text":"Read this exactly."},
966            ])
967        );
968        std::fs::remove_file(path).unwrap();
969    }
970}