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