#![forbid(unsafe_code)]
#![doc = include_str!("../Documentation.md")]
#[cfg(not(unix))]
compile_error!("kcode-k1-codex-websearch requires Unix process groups");
use nix::{
sys::signal::{Signal, killpg},
unistd::Pid,
};
use serde_json::Value;
use std::{ffi::OsString, io, path::PathBuf, time::Instant};
use tokio::{
io::{AsyncRead, AsyncReadExt, AsyncWriteExt},
process::Command,
sync::oneshot,
};
const ARG_ERROR: &str =
"WebSearch failed: model or reasoning_effort is not representable as a process argument";
const PREFLIGHT: &str =
"WebSearch failed: Codex compatibility check failed in kcode-k1-codex-websearch preflight";
const EXEC_IO_ERROR: &str = "WebSearch failed: kcode-k1-codex-websearch could not start or communicate with the Codex execution subprocess";
const EXEC_CANCELLED: &str =
"WebSearch failed: kcode-k1-codex-websearch Codex execution subprocess was cancelled";
const EXEC_EXIT_ERROR: &str =
"WebSearch failed: kcode-k1-codex-websearch Codex execution subprocess exited unsuccessfully";
const EXEC_ERROR: &str = "WebSearch failed: Codex execution failed";
const BAD_JSON: &str = "WebSearch failed: invalid Codex JSONL output";
const INCOMPLETE_ERROR: &str = "WebSearch failed: Codex response was incomplete";
const MESSAGE_ERROR: &str = "WebSearch failed: Codex returned no agent message";
const CONFIG: &str = r#"web_search="live"
tools.web_search=true
tools.view_image=false
apps._default.enabled=false
agents.enabled=false
features.apps=false
features.code_mode.enabled=false
features.goals=false
features.hooks=false
features.memories=false
features.multi_agent=false
features.remote_plugin=false
features.shell_snapshot=false
features.shell_tool=false
features.skill_mcp_dependency_install=false
features.unified_exec=false
memories.generate_memories=false
memories.use_memories=false
history.persistence="none"
check_for_update_on_startup=false
feedback.enabled=false
analytics.enabled=false
allow_login_shell=false
skills.config=[]
mcp_servers={}
plugins={}
marketplaces={}
hooks={}"#;
const EXEC_PREFIX: &str = r#"exec
--search
--ephemeral
--ignore-user-config
--ignore-rules
--json
--sandbox
read-only
--ask-for-approval
never
--skip-git-repo-check
--no-daemon
--strict-config
--model"#;
#[derive(Clone)]
pub struct Runner {
executable: PathBuf,
}
pub struct Request {
pub query: String,
pub model: String,
pub reasoning_effort: String,
pub deadline: Instant,
}
impl Runner {
pub fn new(executable: PathBuf) -> Self {
Self { executable }
}
pub async fn run(&self, request: Request) -> Result<String, String> {
if request.model.contains('\0') || request.reasoning_effort.contains('\0') {
return Err(ARG_ERROR.to_owned());
}
let deadline = request.deadline;
if Instant::now() >= deadline {
return Err(deadline_error("before compatibility preflight"));
}
let mut help = self
.check_help(&["--help"], "codex --help", deadline)
.await?;
help.extend(
self.check_help(&["exec", "--help"], "codex exec --help", deadline)
.await?,
);
let help = String::from_utf8_lossy(&help);
let missing = required_flags()
.into_iter()
.filter(|flag| !help.contains(flag.as_str()))
.collect::<Vec<_>>();
if !missing.is_empty() {
return Err(format!(
"{PREFLIGHT}: combined required-flag validation omitted {}",
missing.join(", ")
));
}
if Instant::now() >= deadline {
return Err(deadline_error("after compatibility preflight"));
}
let output = run_process(
self.executable.clone(),
exec_args(&request.model, &request.reasoning_effort),
request.query.into_bytes(),
deadline,
)
.await
.map_err(|error| match error {
ProcError::Timeout => deadline_error("during Codex execution"),
ProcError::Io => EXEC_IO_ERROR.to_owned(),
ProcError::Cancelled => EXEC_CANCELLED.to_owned(),
})?;
if !output.status.success() {
return Err(EXEC_EXIT_ERROR.to_owned());
}
parse_jsonl(&output.stdout)
}
async fn check_help(
&self,
args: &[&str],
phase: &str,
deadline: Instant,
) -> Result<Vec<u8>, String> {
let output = run_process(
self.executable.clone(),
args.iter().map(OsString::from).collect(),
Vec::new(),
deadline,
)
.await
.map_err(|error| match error {
ProcError::Timeout => {
deadline_error(&format!("during compatibility preflight ({phase})"))
}
ProcError::Io => preflight_error(phase, "subprocess setup, spawn, or I/O failed"),
ProcError::Cancelled => preflight_error(phase, "subprocess was cancelled"),
})?;
if !output.status.success() {
return Err(preflight_error(phase, "subprocess exited unsuccessfully"));
}
let mut text = output.stdout;
text.extend(output.stderr);
Ok(text)
}
}
fn deadline_error(phase: &str) -> String {
format!("WebSearch failed: absolute operation deadline exceeded {phase}")
}
fn preflight_error(phase: &str, cause: &str) -> String {
format!("{PREFLIGHT}: {phase} {cause}")
}
fn exec_args(model: &str, effort: &str) -> Vec<OsString> {
let mut args = EXEC_PREFIX.lines().map(OsString::from).collect::<Vec<_>>();
args.push(model.into());
args.push("-c".into());
args.push(format!("model_reasoning_effort={}", toml_quote(effort)).into());
for config in CONFIG.lines() {
args.push("-c".into());
args.push(config.into());
}
args.push("-".into());
args
}
fn required_flags() -> Vec<String> {
let mut flags = exec_args("", "")
.into_iter()
.filter_map(|arg| arg.into_string().ok())
.filter(|arg| arg.starts_with("--"))
.collect::<Vec<_>>();
flags.push("--config".into());
flags
}
fn toml_quote(value: &str) -> String {
let mut output = String::from("\"");
for character in value.chars() {
match character {
'"' => output.push_str("\\\""),
'\\' => output.push_str("\\\\"),
'\u{8}' => output.push_str("\\b"),
'\t' => output.push_str("\\t"),
'\n' => output.push_str("\\n"),
'\u{c}' => output.push_str("\\f"),
'\r' => output.push_str("\\r"),
value if value.is_control() => output.push_str(&format!("\\u{:04X}", value as u32)),
value => output.push(value),
}
}
output.push('"');
output
}
fn parse_jsonl(bytes: &[u8]) -> Result<String, String> {
let mut turn_completed = false;
let mut turn_failed = false;
let mut message = None;
for raw in bytes.split(|byte| *byte == b'\n') {
let line = raw.trim_ascii();
if line.is_empty() {
continue;
}
let value: Value = serde_json::from_slice(line).map_err(|_| BAD_JSON.to_owned())?;
match value.get("type").and_then(Value::as_str) {
Some("turn.completed") => turn_completed = true,
Some("turn.failed") => turn_failed = true,
Some("item.completed")
if value.pointer("/item/type").and_then(Value::as_str) == Some("agent_message") =>
{
let text = value
.pointer("/item/text")
.and_then(Value::as_str)
.ok_or_else(|| BAD_JSON.to_owned())?;
message = Some(text.to_owned());
}
_ => {}
}
}
if turn_failed {
return Err(EXEC_ERROR.to_owned());
}
if !turn_completed {
return Err(INCOMPLETE_ERROR.to_owned());
}
message.ok_or_else(|| MESSAGE_ERROR.to_owned())
}
enum ProcError {
Io,
Timeout,
Cancelled,
}
async fn drain<R: AsyncRead + Unpin>(mut reader: R) -> io::Result<Vec<u8>> {
let mut bytes = Vec::new();
reader.read_to_end(&mut bytes).await?;
Ok(bytes)
}
async fn run_process(
executable: PathBuf,
args: Vec<OsString>,
input: Vec<u8>,
until: Instant,
) -> Result<std::process::Output, ProcError> {
if Instant::now() >= until {
return Err(ProcError::Timeout);
}
let directory = tempfile::tempdir().map_err(|_| ProcError::Io)?;
let mut command = Command::new(executable);
command
.args(args)
.current_dir(directory.path())
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true)
.process_group(0);
let mut child = command.spawn().map_err(|_| ProcError::Io)?;
let pid = Pid::from_raw(child.id().ok_or(ProcError::Io)? as i32);
let mut stdin = child.stdin.take().ok_or(ProcError::Io)?;
let stdout = child.stdout.take().ok_or(ProcError::Io)?;
let stderr = child.stderr.take().ok_or(ProcError::Io)?;
let (cancel_tx, mut cancel_rx) = oneshot::channel::<()>();
let worker = tokio::spawn(async move {
let _directory = directory;
let complete = async move {
let write = async move {
if input.is_empty() {
return Ok(());
}
stdin.write_all(&input).await?;
stdin.shutdown().await
};
let (status, written, stdout, stderr) =
tokio::join!(child.wait(), write, drain(stdout), drain(stderr));
written.map_err(|_| ProcError::Io)?;
Ok::<_, ProcError>(std::process::Output {
status: status.map_err(|_| ProcError::Io)?,
stdout: stdout.map_err(|_| ProcError::Io)?,
stderr: stderr.map_err(|_| ProcError::Io)?,
})
};
tokio::pin!(complete);
let timer = tokio::time::sleep_until(tokio::time::Instant::from_std(until));
tokio::pin!(timer);
let error = tokio::select! {
result = &mut complete => return result,
_ = &mut timer => ProcError::Timeout,
_ = &mut cancel_rx => ProcError::Cancelled,
};
let _ = killpg(pid, Signal::SIGKILL);
let _ = (&mut complete).await;
Err(error)
});
let result = worker.await.map_err(|_| ProcError::Io)?;
drop(cancel_tx);
result
}
#[cfg(test)]
mod tests {
use super::*;
use std::{fs, os::unix::fs::PermissionsExt, path::Path, time::Duration};
struct Fake(tempfile::TempDir);
impl Fake {
fn file(&self, name: &str) -> PathBuf {
self.0.path().join(name)
}
async fn run(&self, duration: Duration) -> Result<String, String> {
Runner::new(self.file("codex"))
.run(request("q", duration))
.await
}
}
fn fake(body: &str, exit: i32, help: Option<String>) -> Fake {
let directory = tempfile::tempdir().unwrap();
let executable = directory.path().join("codex");
let flags = required_flags().join(" ");
let help = help.unwrap_or_else(|| format!("printf '%s\\n' '{flags}'"));
let script = format!(
"#!/bin/sh\nbase=${{0%/*}}\nif [ \"$1\" = \"--help\" ] || {{ [ \"$1\" = exec ] && [ \"$2\" = \"--help\" ]; }}; then\n{}\nexit 0\nfi\nprintf '%s\\n' \"$@\" > \"$base/args\"\ncat > \"$base/input\"\n{}\nexit {}\n",
help, body, exit,
);
fs::write(&executable, script).unwrap();
fs::set_permissions(&executable, fs::Permissions::from_mode(0o755)).unwrap();
Fake(directory)
}
fn request(query: &str, duration: Duration) -> Request {
Request {
query: query.into(),
model: "model-x".into(),
reasoning_effort: "high".into(),
deadline: Instant::now() + duration,
}
}
const SUCCESS: &str = r#"printf '%s\n' \
'{"type":"item.completed","item":{"type":"agent_message","text":"answer"}}' \
'{"type":"turn.completed"}'"#;
const DESCENDANTS: &str = r#"sleep 30 & echo "$$ $!" > "$base/pids"; wait"#;
#[tokio::test]
async fn invocation_stdin_and_last_message_are_exact() {
let fake = fake(
r#"printf '%s\n' \
'{"type":"future.event","x":1}' \
'{"type":"item.completed","item":{"type":"agent_message","text":"old"}}' \
'{"type":"item.completed","item":{"type":"agent_message","text":"最後\nline"}}' \
'{"type":"turn.completed","extra":true}'"#,
0,
None,
);
let query = "héllo\n世界\0tail";
let answer = Runner::new(fake.file("codex"))
.run(request(query, Duration::from_secs(5)))
.await
.unwrap();
assert_eq!(answer, "最後\nline");
assert_eq!(fs::read(fake.file("input")).unwrap(), query.as_bytes());
let text = fs::read_to_string(fake.file("args")).unwrap();
let got = text.lines().collect::<Vec<_>>();
assert_eq!(&got[..3], ["exec", "--search", "--ephemeral"]);
assert_eq!(
&got[3..6],
["--ignore-user-config", "--ignore-rules", "--json"]
);
assert_eq!(
&got[6..10],
["--sandbox", "read-only", "--ask-for-approval", "never"]
);
assert_eq!(
&got[10..13],
["--skip-git-repo-check", "--no-daemon", "--strict-config"]
);
assert_eq!(&got[13..15], ["--model", "model-x"]);
assert_eq!(&got[15..17], ["-c", r#"model_reasoning_effort="high""#]);
for (index, config) in CONFIG.lines().enumerate() {
assert_eq!(&got[17 + index * 2..19 + index * 2], ["-c", config]);
}
assert_eq!(got.last(), Some(&"-"));
assert_eq!(toml_quote("a\"\n\u{7f}"), "\"a\\\"\\n\\u007F\"");
}
#[tokio::test]
async fn compatibility_diagnostic_names_missing_flag_without_help() {
let advertised = required_flags().join(" ").replace("--no-daemon", "");
let help = format!("printf '%s\\n' '{advertised} conspicuous-private-help-text'");
let fake = fake("", 0, Some(help));
let error = fake.run(Duration::from_secs(5)).await.unwrap_err();
assert_eq!(
error,
preflight_error("combined required-flag validation", "omitted --no-daemon")
);
assert!(!error.contains("conspicuous-private-help-text"));
}
#[tokio::test]
async fn help_and_execution_use_only_the_absolute_deadline() {
let help = format!(
"if [ \"$1\" = \"--help\" ]; then sleep 3; fi\nprintf '%s\\n' '{}'",
required_flags().join(" ")
);
let slow_help = fake(SUCCESS, 0, Some(help));
assert_eq!(
slow_help.run(Duration::from_secs(7)).await.unwrap(),
"answer"
);
let body = format!("sleep 2\n{SUCCESS}");
let slow_execution = fake(&body, 0, None);
assert_eq!(
slow_execution
.run(Duration::from_millis(2800))
.await
.unwrap(),
"answer"
);
}
#[tokio::test]
async fn rejects_failures_and_nul_arguments() {
let cases = [
("{", 0, BAD_JSON),
(r#"{"type":"turn.failed"}"#, 0, EXEC_ERROR),
(r#"{"type":"turn.completed"}"#, 7, EXEC_EXIT_ERROR),
(r#"{"type":"turn.completed"}"#, 0, MESSAGE_ERROR),
(
r#"{"type":"item.completed","item":{"type":"agent_message","text":"x"}}"#,
0,
INCOMPLETE_ERROR,
),
];
for (json, status, error) in cases {
let body = format!("printf '%s\\n' '{json}'");
let fake = fake(&body, status, None);
assert_eq!(fake.run(Duration::from_secs(5)).await.unwrap_err(), error);
}
let runner = Runner::new("/not/executed".into());
for (model, effort) in [("bad\0", "high"), ("model-x", "bad\0")] {
let mut bad = request("q", Duration::from_secs(2));
bad.model = model.into();
bad.reasoning_effort = effort.into();
assert_eq!(runner.run(bad).await.unwrap_err(), ARG_ERROR);
}
}
async fn wait_until(mut condition: impl FnMut() -> bool) {
for _ in 0..200 {
if condition() {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("condition was not reached");
}
fn process_is_live(pid: i32) -> bool {
fs::read_to_string(format!("/proc/{pid}/stat")).is_ok_and(|stat| {
let state = stat
.rsplit_once(") ")
.and_then(|(_, tail)| tail.bytes().next());
!matches!(state, Some(b'Z' | b'X'))
})
}
async fn wait_for_processes_to_stop(path: &Path) {
let pids = fs::read_to_string(path)
.unwrap()
.split_whitespace()
.map(|value| value.parse::<i32>().unwrap())
.collect::<Vec<_>>();
assert!(pids.len() >= 2);
wait_until(|| pids.iter().all(|pid| !process_is_live(*pid))).await;
}
#[tokio::test]
async fn deadline_and_cancellation_reap_descendants() {
let deadline = fake(DESCENDANTS, 0, None);
let error = deadline.run(Duration::from_millis(1600)).await.unwrap_err();
assert_eq!(error, deadline_error("during Codex execution"));
wait_for_processes_to_stop(&deadline.file("pids")).await;
let cancelled = fake(DESCENDANTS, 0, None);
let runner = Runner::new(cancelled.file("codex"));
let task =
tokio::spawn(async move { runner.run(request("q", Duration::from_secs(10))).await });
wait_until(|| cancelled.file("pids").exists()).await;
task.abort();
let _ = task.await;
wait_for_processes_to_stop(&cancelled.file("pids")).await;
}
}