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    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}