kcode-k1-codex-websearch 0.1.2

Isolated one-shot Codex live-Web search runner
Documentation
#![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
}