#![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, os::unix::process::ExitStatusExt, 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_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 STDERR_RULES: &[(&[&str], &str)] = &[
(
&[
"unexpected argument",
"unrecognized argument",
"unknown argument",
"unexpected option",
"unrecognized option",
"unknown option",
"configuration",
"config key",
"config value",
"strict config",
],
"CLI argument/configuration",
),
(
&[
"authentication",
"not logged in",
"login required",
"credential",
],
"authentication",
),
(&["unsupported model", "unknown model"], "unsupported model"),
(&["rate limit", "quota"], "quota/rate limit"),
(
&["network", "connection", "dns", "transport"],
"network/transport",
),
(
&["upstream", "provider", "service unavailable"],
"provider/upstream",
),
(&["codex failed", "codex error"], "other Codex failure"),
];
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(execution_exit_error(output.status, &output.stderr));
}
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 execution_exit_error(status: std::process::ExitStatus, stderr: &[u8]) -> String {
let termination = status.code().map_or_else(
|| {
format!(
"terminated by signal {}",
status.signal().unwrap_or_default()
)
},
|code| format!("exited with numeric exit code {code}"),
);
format!(
"WebSearch failed: kcode-k1-codex-websearch Codex execution phase {termination}; stderr classification: {}",
classify_stderr(stderr)
)
}
fn classify_stderr(bytes: &[u8]) -> &'static str {
if bytes.is_empty() {
return "empty";
}
let Ok(text) = std::str::from_utf8(bytes) else {
return "non-UTF8";
};
let text = text.to_ascii_lowercase();
STDERR_RULES
.iter()
.find_map(|(needles, category)| {
needles
.iter()
.any(|needle| text.contains(needle))
.then_some(*category)
})
.unwrap_or("unclassified")
}
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
}