use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::process::{Child, Command};
use tokio::sync::Mutex;
use tokio_util::sync::CancellationToken;
use bamboo_agent_core::{AgentEvent, TokenUsage, ToolResult};
use bamboo_subagent::codex_discovery::discover_codex_cli;
use bamboo_subagent::executor::{ChildExecutor, ChildOutcome, EventSink, SteerInbox};
use bamboo_subagent::executor_util::{build_rehydrated_turn, write_json_atomic};
use bamboo_subagent::proto::RunSpec;
pub use bamboo_subagent::codex_discovery::MIN_CODEX_VERSION;
const MAX_STDOUT_LINE_BYTES: usize = 10 * 1024 * 1024;
const STDERR_TAIL_BYTES: usize = 16 * 1024;
const TOOL_RESULT_TRUNCATE_CHARS: usize = 20_000;
const SIGTERM_WAIT: Duration = Duration::from_secs(5);
const PROCESS_EXIT_WAIT: Duration = Duration::from_secs(5);
const ENV_ALLOWLIST: &[&str] = &[
"HOME", "PATH", "SHELL", "TERM", "LANG", "TMPDIR", "USER", "LOGNAME",
];
const CODEX_PROVIDER_ENV: &str = "BAMBOO_CODEX_PROVIDER_KEY";
const CODEX_SESSION_STATE_FILE: &str = "codex-session.json";
#[derive(Debug, Clone, Serialize, Deserialize)]
struct CodexSessionState {
thread_id: String,
workspace: Option<String>,
codex_home_mode: String,
updated_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CodexAuthMode {
Inherit,
ApiKey,
Custom,
Bamboo,
}
impl CodexAuthMode {
pub(crate) fn as_str(self) -> &'static str {
match self {
Self::Inherit => "inherit",
Self::ApiKey => "api_key",
Self::Custom => "custom",
Self::Bamboo => "bamboo",
}
}
}
#[derive(Clone)]
pub struct CodexAuthConfig {
mode: CodexAuthMode,
base_url: Option<String>,
wire_api: String,
provider_key: Option<String>,
}
impl CodexAuthConfig {
pub fn inherit() -> Self {
Self {
mode: CodexAuthMode::Inherit,
base_url: None,
wire_api: "responses".to_string(),
provider_key: None,
}
}
pub fn mode(&self) -> CodexAuthMode {
self.mode
}
pub(crate) fn isolated(&self) -> bool {
self.mode != CodexAuthMode::Inherit
}
fn generated_config_toml(&self) -> Result<String, String> {
if !self.isolated() || self.mode == CodexAuthMode::ApiKey {
return Ok("# Generated by Bamboo; authentication is environment-only.\n".to_string());
}
let provider_id = match self.mode {
CodexAuthMode::Custom => "custom",
CodexAuthMode::Bamboo => "bamboo",
CodexAuthMode::Inherit | CodexAuthMode::ApiKey => unreachable!(),
};
let base_url = self.base_url.as_ref().ok_or_else(|| {
format!("Codex auth mode '{provider_id}' requires a provider base URL")
})?;
let mut provider = toml::Table::new();
provider.insert(
"name".to_string(),
toml::Value::String(format!("Bamboo {provider_id}")),
);
provider.insert(
"base_url".to_string(),
toml::Value::String(base_url.clone()),
);
provider.insert(
"env_key".to_string(),
toml::Value::String(CODEX_PROVIDER_ENV.to_string()),
);
provider.insert(
"wire_api".to_string(),
toml::Value::String(self.wire_api.clone()),
);
let mut providers = toml::Table::new();
providers.insert(provider_id.to_string(), toml::Value::Table(provider));
let mut root = toml::Table::new();
root.insert(
"model_provider".to_string(),
toml::Value::String(provider_id.to_string()),
);
root.insert("model_providers".to_string(), toml::Value::Table(providers));
toml::to_string(&toml::Value::Table(root))
.map_err(|error| format!("serialize isolated Codex config.toml: {error}"))
}
pub(crate) fn generated_app_server_config_toml(
&self,
token_helper: &Path,
token_path: &Path,
) -> Result<String, String> {
if self.mode != CodexAuthMode::Bamboo {
return self.generated_config_toml();
}
let base_url = self
.base_url
.as_ref()
.ok_or_else(|| "Codex auth mode 'bamboo' requires a provider base URL".to_string())?;
let mut auth = toml::Table::new();
auth.insert(
"command".to_string(),
toml::Value::String(token_helper.to_string_lossy().into_owned()),
);
auth.insert(
"args".to_string(),
toml::Value::Array(vec![
toml::Value::String("codex-provider-token".to_string()),
toml::Value::String(token_path.to_string_lossy().into_owned()),
]),
);
auth.insert("timeout_ms".to_string(), toml::Value::Integer(5_000));
auth.insert("refresh_interval_ms".to_string(), toml::Value::Integer(1));
let mut provider = toml::Table::new();
provider.insert(
"name".to_string(),
toml::Value::String("Bamboo bamboo".to_string()),
);
provider.insert(
"base_url".to_string(),
toml::Value::String(base_url.clone()),
);
provider.insert(
"wire_api".to_string(),
toml::Value::String(self.wire_api.clone()),
);
provider.insert("auth".to_string(), toml::Value::Table(auth));
let mut providers = toml::Table::new();
providers.insert("bamboo".to_string(), toml::Value::Table(provider));
let mut root = toml::Table::new();
root.insert(
"model_provider".to_string(),
toml::Value::String("bamboo".to_string()),
);
root.insert("model_providers".to_string(), toml::Value::Table(providers));
toml::to_string(&toml::Value::Table(root))
.map_err(|error| format!("serialize isolated Codex app-server config.toml: {error}"))
}
pub(crate) fn provider_key(&self) -> Option<&str> {
self.provider_key.as_deref()
}
}
pub fn read_codex_provider_token(path: &Path) -> Result<String, String> {
#[cfg(not(unix))]
{
let metadata = std::fs::symlink_metadata(path)
.map_err(|error| format!("inspect Codex provider token: {error}"))?;
if metadata.file_type().is_symlink() {
return Err("Codex provider token path must not be a symlink".to_string());
}
}
let mut options = std::fs::OpenOptions::new();
options.read(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC);
}
let mut file = options
.open(path)
.map_err(|error| format!("open Codex provider token: {error}"))?;
let metadata = file
.metadata()
.map_err(|error| format!("inspect Codex provider token: {error}"))?;
if !metadata.is_file() {
return Err("Codex provider token path must be a regular file".to_string());
}
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt as _;
if metadata.permissions().mode() & 0o077 != 0 {
return Err(
"Codex provider token file must not be accessible by group/other".to_string(),
);
}
}
let mut token = String::new();
std::io::Read::read_to_string(&mut file, &mut token)
.map_err(|error| format!("read Codex provider token: {error}"))?;
let token = token.trim();
if token.is_empty() {
return Err("Codex provider token file is empty".to_string());
}
Ok(token.to_string())
}
#[allow(clippy::too_many_arguments)]
pub fn resolve_codex_auth_config(
auth_mode: Option<&str>,
legacy_inherit_user_config: bool,
base_url: Option<String>,
wire_api: Option<String>,
provider_key_ref: Option<&str>,
credentials: &[bamboo_subagent::provision::ScopedCredential],
forward_env: &[String],
) -> Result<CodexAuthConfig, String> {
let mode = match auth_mode {
Some("inherit") => CodexAuthMode::Inherit,
Some("api_key") => CodexAuthMode::ApiKey,
Some("custom") => CodexAuthMode::Custom,
Some("bamboo") => CodexAuthMode::Bamboo,
Some(other) => {
return Err(format!(
"unknown Codex auth mode '{other}'; expected inherit, api_key, custom, or bamboo"
))
}
None if legacy_inherit_user_config => CodexAuthMode::Inherit,
None => CodexAuthMode::Bamboo,
};
let wire_api = wire_api.unwrap_or_else(|| "responses".to_string());
if wire_api != "responses" {
return Err(format!(
"unsupported Codex wire_api '{wire_api}'; Codex CLI >= 0.144 requires responses"
));
}
validate_codex_base_url(mode, base_url.as_deref())?;
validate_codex_forward_env(mode, forward_env)?;
let provider_key = if mode == CodexAuthMode::Custom {
let reference = provider_key_ref.ok_or_else(|| {
"Codex auth mode 'custom' requires codex_provider_key_ref".to_string()
})?;
Some(
credentials
.iter()
.find(|credential| credential.credential_ref.as_deref() == Some(reference))
.map(|credential| credential.api_key.clone())
.ok_or_else(|| {
format!(
"Codex custom provider credential reference '{reference}' did not resolve"
)
})?,
)
} else {
None
};
Ok(CodexAuthConfig {
mode,
base_url,
wire_api,
provider_key,
})
}
pub fn resolve_codex_state_dir(storage_dir: &Option<String>, child_id: &str) -> PathBuf {
storage_dir
.clone()
.map(PathBuf::from)
.unwrap_or_else(|| bamboo_config::paths::subagents_dir().join(child_id))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CodexSandbox {
ReadOnly,
WorkspaceWrite,
DangerFullAccess,
}
impl CodexSandbox {
fn parse(raw: &str) -> Result<Self, String> {
match raw {
"read-only" => Ok(Self::ReadOnly),
"workspace-write" => Ok(Self::WorkspaceWrite),
"danger-full-access" => Ok(Self::DangerFullAccess),
other => Err(format!(
"unknown Codex sandbox '{other}'; expected read-only, workspace-write, or danger-full-access"
)),
}
}
fn as_str(self) -> &'static str {
match self {
Self::ReadOnly => "read-only",
Self::WorkspaceWrite => "workspace-write",
Self::DangerFullAccess => "danger-full-access",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CodexApprovalPolicy {
Never,
OnFailure,
OnRequest,
}
impl CodexApprovalPolicy {
fn parse(raw: &str) -> Result<Self, String> {
match raw {
"never" => Ok(Self::Never),
"on-failure" => Ok(Self::OnFailure),
"on-request" => Ok(Self::OnRequest),
"untrusted" => Err("Codex approval policy 'untrusted' is unsupported; use on-request in app_server mode".to_string()),
other => Err(format!(
"unknown Codex approval policy '{other}'; expected never, on-failure, or on-request"
)),
}
}
fn as_str(self) -> &'static str {
match self {
Self::Never => "never",
Self::OnFailure => "on-failure",
Self::OnRequest => "on-request",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CodexPolicyInvocation {
Explicit,
FullAuto,
DangerBypass,
}
impl CodexPolicyInvocation {
fn as_str(self) -> &'static str {
match self {
Self::Explicit => "explicit",
Self::FullAuto => "full-auto",
Self::DangerBypass => "danger-bypass",
}
}
}
#[derive(Debug, Clone)]
pub struct CodexPermissionConfig {
sandbox: Option<CodexSandbox>,
approval_policy: Option<CodexApprovalPolicy>,
network_access: bool,
allow_danger_bypass: bool,
permission_profile: Option<String>,
provisioned_bypass: bool,
workspace_owned: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct EffectiveCodexPolicy {
sandbox: CodexSandbox,
approval_policy: CodexApprovalPolicy,
invocation: CodexPolicyInvocation,
network_access: bool,
warnings: Vec<String>,
}
impl CodexPermissionConfig {
fn effective(&self, bypass: bool, is_root: bool) -> EffectiveCodexPolicy {
let read_only_profile = self
.permission_profile
.as_deref()
.is_some_and(profile_is_read_only);
let mut warnings = Vec::new();
let (sandbox, approval_policy, invocation) = match self.sandbox {
Some(CodexSandbox::ReadOnly) => (
CodexSandbox::ReadOnly,
self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
CodexPolicyInvocation::Explicit,
),
Some(CodexSandbox::WorkspaceWrite) => (
CodexSandbox::WorkspaceWrite,
self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
CodexPolicyInvocation::Explicit,
),
Some(CodexSandbox::DangerFullAccess) => {
resolve_danger_request(bypass, self.allow_danger_bypass, is_root, &mut warnings)
}
None if self.allow_danger_bypass => {
resolve_danger_request(bypass, true, is_root, &mut warnings)
}
None if bypass => (
CodexSandbox::WorkspaceWrite,
CodexApprovalPolicy::Never,
CodexPolicyInvocation::FullAuto,
),
None if read_only_profile => (
CodexSandbox::ReadOnly,
self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
CodexPolicyInvocation::Explicit,
),
None => (
CodexSandbox::WorkspaceWrite,
self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
CodexPolicyInvocation::Explicit,
),
};
let network_access = match sandbox {
CodexSandbox::ReadOnly => false,
CodexSandbox::WorkspaceWrite => self.network_access,
CodexSandbox::DangerFullAccess => true,
};
if invocation == CodexPolicyInvocation::DangerBypass {
warnings.push(
"DANGER: Codex is running with approvals and the OS sandbox disabled for this parent-bypass run"
.to_string(),
);
}
EffectiveCodexPolicy {
sandbox,
approval_policy,
invocation,
network_access,
warnings,
}
}
}
fn resolve_danger_request(
bypass: bool,
allow_danger_bypass: bool,
is_root: bool,
warnings: &mut Vec<String>,
) -> (CodexSandbox, CodexApprovalPolicy, CodexPolicyInvocation) {
if bypass && allow_danger_bypass && !is_root {
return (
CodexSandbox::DangerFullAccess,
CodexApprovalPolicy::Never,
CodexPolicyInvocation::DangerBypass,
);
}
let reason = if !bypass {
"the parent run is not in bypass mode"
} else if !allow_danger_bypass {
"codex_allow_danger_bypass is false"
} else {
"the worker is running as root"
};
warnings.push(format!(
"Codex danger-full-access request was downgraded to a sandboxed policy because {reason}"
));
(
CodexSandbox::WorkspaceWrite,
CodexApprovalPolicy::Never,
CodexPolicyInvocation::FullAuto,
)
}
fn profile_is_read_only(raw: &str) -> bool {
let normalized = raw.trim().to_ascii_lowercase().replace(['_', ' '], "-");
matches!(
normalized.as_str(),
"read-only" | "readonly" | "research" | "researcher" | "guardian" | "plan"
)
}
#[allow(clippy::too_many_arguments)]
pub fn resolve_codex_permission_config(
sandbox: Option<&str>,
approval_policy: Option<&str>,
network_access: bool,
allow_danger_bypass: bool,
permission_profile: Option<String>,
provisioned_bypass: bool,
workspace_owned: bool,
) -> Result<CodexPermissionConfig, String> {
let sandbox = sandbox.map(CodexSandbox::parse).transpose()?;
let approval_policy = approval_policy
.map(CodexApprovalPolicy::parse)
.transpose()?;
if approval_policy == Some(CodexApprovalPolicy::OnRequest) {
return Err(
"Codex approval policy 'on-request' requires codex_mode = \"app_server\"; non-interactive exec mode has no approval relay"
.to_string(),
);
}
let profile_read_only = permission_profile
.as_deref()
.is_some_and(profile_is_read_only);
if network_access
&& (sandbox == Some(CodexSandbox::ReadOnly) || (sandbox.is_none() && profile_read_only))
{
return Err(
"Codex network access requires an effective workspace-write sandbox".to_string(),
);
}
Ok(CodexPermissionConfig {
sandbox,
approval_policy,
network_access,
allow_danger_bypass,
permission_profile,
provisioned_bypass,
workspace_owned,
})
}
#[allow(clippy::too_many_arguments)]
pub fn resolve_codex_app_server_permission_config(
sandbox: Option<&str>,
approval_policy: Option<&str>,
network_access: bool,
allow_danger_bypass: bool,
permission_profile: Option<String>,
provisioned_bypass: bool,
workspace_owned: bool,
) -> Result<CodexPermissionConfig, String> {
match approval_policy {
None | Some("on-request") => {}
Some(other) => {
return Err(format!(
"Codex approval policy '{other}' is incompatible with codex_mode = \"app_server\"; use on-request"
))
}
}
let sandbox = sandbox.map(CodexSandbox::parse).transpose()?;
let profile_read_only = permission_profile
.as_deref()
.is_some_and(profile_is_read_only);
if network_access
&& (sandbox == Some(CodexSandbox::ReadOnly) || (sandbox.is_none() && profile_read_only))
{
return Err(
"Codex network access requires an effective workspace-write sandbox".to_string(),
);
}
Ok(CodexPermissionConfig {
sandbox,
approval_policy: Some(CodexApprovalPolicy::OnRequest),
network_access,
allow_danger_bypass,
permission_profile,
provisioned_bypass,
workspace_owned,
})
}
impl CodexPermissionConfig {
pub(crate) fn app_server_posture(&self, bypass: bool) -> (String, bool, Vec<String>) {
let mut policy = self.effective(bypass, running_as_root());
policy.approval_policy = CodexApprovalPolicy::OnRequest;
(
policy.sandbox.as_str().to_string(),
policy.network_access,
policy.warnings,
)
}
pub(crate) fn permission_profile(&self) -> Option<&str> {
self.permission_profile.as_deref()
}
pub(crate) fn provisioned_bypass(&self) -> bool {
self.provisioned_bypass
}
}
#[cfg(unix)]
fn running_as_root() -> bool {
unsafe { libc::geteuid() == 0 }
}
#[cfg(not(unix))]
fn running_as_root() -> bool {
false
}
pub struct CodexExecutor {
binary: PathBuf,
version: String,
model: Option<String>,
permissions: CodexPermissionConfig,
workspace: Option<String>,
state_dir: Option<PathBuf>,
forward_env: Vec<String>,
auth: CodexAuthConfig,
}
impl CodexExecutor {
#[allow(clippy::too_many_arguments)]
pub async fn new(
binary: Option<String>,
model: Option<String>,
workspace: Option<String>,
state_dir: Option<PathBuf>,
forward_env: Vec<String>,
auth: CodexAuthConfig,
permissions: CodexPermissionConfig,
) -> Result<Self, String> {
let discovery = discover_codex_cli(binary.as_deref()).await?;
let executor = Self {
binary: PathBuf::from(discovery.path),
version: discovery.version,
model,
permissions,
workspace,
state_dir,
forward_env,
auth,
};
executor.prepare_auth_home().await?;
Ok(executor)
}
fn last_message_path(&self) -> Option<PathBuf> {
self.state_dir
.as_ref()
.map(|directory| directory.join("codex-last-message.txt"))
}
fn session_state_path(&self) -> Option<PathBuf> {
self.state_dir
.as_ref()
.map(|directory| directory.join(CODEX_SESSION_STATE_FILE))
}
fn codex_home_mode(&self) -> &'static str {
if self.auth.isolated() {
"isolated"
} else {
"inherit"
}
}
async fn read_session_state(&self) -> Option<CodexSessionState> {
let path = self.session_state_path()?;
let bytes = tokio::fs::read(path).await.ok()?;
serde_json::from_slice(&bytes).ok()
}
async fn delete_session_state(&self) {
let Some(path) = self.session_state_path() else {
return;
};
match tokio::fs::remove_file(&path).await {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => tracing::warn!(
path = %path.display(),
%error,
"codex: remove stale session state"
),
}
}
async fn write_session_state(&self, thread_id: &str) {
let Some(path) = self.session_state_path() else {
return;
};
let state = CodexSessionState {
thread_id: thread_id.to_string(),
workspace: self.workspace.clone(),
codex_home_mode: self.codex_home_mode().to_string(),
updated_at: Utc::now(),
};
if let Err(error) = write_json_atomic(&path, &state).await {
tracing::warn!(%error, "codex: persist thread state");
}
}
async fn resolve_resume_id(&self) -> Option<String> {
let state = self.read_session_state().await?;
if state.workspace != self.workspace {
tracing::warn!(
recorded = ?state.workspace,
current = ?self.workspace,
"codex: session state workspace mismatch; using history rehydration"
);
return None;
}
let current_mode = self.codex_home_mode();
if state.codex_home_mode != current_mode {
tracing::warn!(
recorded = %state.codex_home_mode,
current = current_mode,
"codex: session state CODEX_HOME mode mismatch; using history rehydration"
);
return None;
}
(!state.thread_id.trim().is_empty()).then_some(state.thread_id)
}
fn build_command(
&self,
run_provider_token: Option<&str>,
policy: &EffectiveCodexPolicy,
resume_id: Option<&str>,
) -> Result<Command, String> {
let mut command = Command::new(&self.binary);
command
.arg("exec")
.arg("--json")
.arg("--color")
.arg("never");
let workspace_path = if let Some(workspace) = &self.workspace {
command.arg("--cd").arg(workspace);
Some(PathBuf::from(workspace))
} else {
std::env::current_dir().ok()
};
if self.permissions.workspace_owned
&& workspace_path
.as_deref()
.is_some_and(|path| !has_git_metadata(path))
{
command.arg("--skip-git-repo-check");
}
if let Some(model) = &self.model {
command.arg("--model").arg(model);
}
match policy.invocation {
CodexPolicyInvocation::Explicit => {
command.arg("--sandbox").arg(policy.sandbox.as_str());
command.arg("--config").arg(format!(
"approval_policy=\"{}\"",
policy.approval_policy.as_str()
));
}
CodexPolicyInvocation::FullAuto => {
command.arg("--full-auto");
}
CodexPolicyInvocation::DangerBypass => {
command.arg("--dangerously-bypass-approvals-and-sandbox");
}
}
if policy.sandbox == CodexSandbox::WorkspaceWrite && policy.network_access {
command
.arg("--config")
.arg("sandbox_workspace_write.network_access=true");
}
if self.auth.isolated() {
command.arg("--ignore-rules");
}
if let Some(path) = self.last_message_path() {
command.arg("--output-last-message").arg(path);
}
if let Some(thread_id) = resume_id {
command.arg("resume").arg(thread_id);
}
command.arg("-");
command.env_clear();
for (key, value) in std::env::vars() {
if ENV_ALLOWLIST.contains(&key.as_str()) || key.starts_with("LC_") {
command.env(key, value);
}
}
if let Some(home) = self.codex_home() {
command.env("CODEX_HOME", home);
}
for name in &self.forward_env {
if let Ok(value) = std::env::var(name) {
command.env(name, value);
}
}
match self.auth.mode {
CodexAuthMode::Custom => {
let key = self.auth.provider_key.as_deref().ok_or_else(|| {
"Codex custom provider key was not resolved at provisioning".to_string()
})?;
command.env(CODEX_PROVIDER_ENV, key);
}
CodexAuthMode::Bamboo => {
let token = run_provider_token.ok_or_else(|| {
"Codex bamboo auth requires a per-run provider token".to_string()
})?;
command.env(CODEX_PROVIDER_ENV, token);
}
CodexAuthMode::Inherit | CodexAuthMode::ApiKey => {}
}
command.stdin(Stdio::piped());
command.stdout(Stdio::piped());
command.stderr(Stdio::piped());
command.kill_on_drop(true);
#[cfg(unix)]
command.process_group(0);
Ok(command)
}
fn codex_home(&self) -> Option<PathBuf> {
self.auth
.isolated()
.then(|| self.state_dir.as_ref().map(|path| path.join("codex-home")))
.flatten()
}
async fn prepare_auth_home(&self) -> Result<(), String> {
if !self.auth.isolated() {
return Ok(());
}
let home = self.codex_home().ok_or_else(|| {
"isolated Codex auth requires a Bamboo-managed state directory".to_string()
})?;
tokio::fs::create_dir_all(&home)
.await
.map_err(|error| format!("create isolated CODEX_HOME '{}': {error}", home.display()))?;
#[cfg(unix)]
tokio::fs::set_permissions(&home, std::os::unix::fs::PermissionsExt::from_mode(0o700))
.await
.map_err(|error| format!("secure isolated CODEX_HOME '{}': {error}", home.display()))?;
let auth_path = home.join("auth.json");
match tokio::fs::remove_file(&auth_path).await {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => {
return Err(format!(
"remove stale isolated Codex auth '{}': {error}",
auth_path.display()
))
}
}
let config_path = home.join("config.toml");
tokio::fs::write(&config_path, self.auth.generated_config_toml()?)
.await
.map_err(|error| {
format!(
"write isolated Codex config '{}': {error}",
config_path.display()
)
})?;
#[cfg(unix)]
tokio::fs::set_permissions(
&config_path,
std::os::unix::fs::PermissionsExt::from_mode(0o600),
)
.await
.map_err(|error| {
format!(
"secure isolated Codex config '{}': {error}",
config_path.display()
)
})?;
Ok(())
}
async fn prepare_output_file(&self) -> Result<(), String> {
let Some(path) = self.last_message_path() else {
return Ok(());
};
if let Some(parent) = path.parent() {
tokio::fs::create_dir_all(parent).await.map_err(|error| {
format!("create Codex state dir '{}': {error}", parent.display())
})?;
}
match tokio::fs::remove_file(&path).await {
Ok(()) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(format!(
"remove stale Codex last-message file '{}': {error}",
path.display()
)),
}
}
fn handle_event(
&self,
value: Value,
policy: &EffectiveCodexPolicy,
events: &EventSink,
state: &mut RunState,
) -> Option<String> {
let event_type = value.get("type").and_then(Value::as_str).unwrap_or("");
let mut captured_thread_id = None;
match event_type {
"thread.started" => {
state.thread_id = value
.get("thread_id")
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
if !state.thread_id.is_empty() {
captured_thread_id = Some(state.thread_id.clone());
}
events.emit(json!({
"type": "runner_progress",
"session_id": state.thread_id,
"round_count": 0,
"executor": "codex",
"binary": self.binary,
"version": self.version,
"model": self.model,
"auth_mode": self.auth.mode().as_str(),
"codex_home_mode": self.codex_home_mode(),
"forward_env": self.forward_env,
"sandbox": policy.sandbox.as_str(),
"approval_policy": policy.approval_policy.as_str(),
"network_access": policy.network_access,
"policy_invocation": policy.invocation.as_str(),
"permission_profile": self.permissions.permission_profile.as_deref(),
}));
}
"turn.started" => {
state.turn_started = true;
events.emit(event_json(AgentEvent::RunnerProgress {
session_id: state.session_id(),
round_count: 1,
}));
}
"item.started" | "item.updated" | "item.completed" => {
let phase = event_type.trim_start_matches("item.");
if let Some(item) = value.get("item") {
handle_item(phase, item, events, state);
if phase == "completed"
&& policy.sandbox != CodexSandbox::DangerFullAccess
&& !state.tool_error_emitted
&& item.get("type").and_then(Value::as_str) == Some("agent_message")
&& item
.get("text")
.and_then(Value::as_str)
.is_some_and(looks_like_sandbox_denial)
{
let item_id = item
.get("id")
.and_then(Value::as_str)
.unwrap_or("codex-sandbox-denial");
let error = item.get("text").and_then(Value::as_str).unwrap_or(
"Codex reported that the OS sandbox denied a tool operation",
);
events.emit(event_json(AgentEvent::ToolError {
tool_call_id: format!("{item_id}-sandbox-denial"),
error: truncate_chars(error, TOOL_RESULT_TRUNCATE_CHARS),
}));
state.tool_error_emitted = true;
}
}
}
"turn.completed" => {
state.completed = true;
state.usage = parse_usage(value.get("usage"));
if let Some(text) = final_text_from_terminal(&value) {
state.last_agent_message = text;
}
events.emit(event_json(AgentEvent::Complete { usage: state.usage }));
}
"turn.failed" => {
let message = error_message(&value, "Codex turn failed");
state.failure = Some(message.clone());
events.emit(event_json(AgentEvent::Error { message }));
}
"error" => {
let message = error_message(&value, "Codex CLI error");
state.failure = Some(message.clone());
events.emit(event_json(AgentEvent::Error { message }));
}
other => {
tracing::debug!(event_type = other, "codex: unrecognized JSONL event");
}
}
captured_thread_id
}
async fn read_last_message(&self) -> Option<String> {
let path = self.last_message_path()?;
let text = tokio::fs::read_to_string(path).await.ok()?;
let trimmed = text.trim();
(!trimmed.is_empty()).then(|| trimmed.to_string())
}
async fn run_process(
&self,
prompt: &str,
resume_id: Option<&str>,
run_provider_token: Option<&str>,
parent_bypass: bool,
events: &EventSink,
cancel: &CancellationToken,
) -> (ChildOutcome, bool) {
let policy = self.permissions.effective(parent_bypass, running_as_root());
self.emit_policy_bootstrap(&policy, events);
if let Err(error) = self.prepare_auth_home().await {
return (ChildOutcome::error(error), false);
}
if let Err(error) = self.prepare_output_file().await {
return (ChildOutcome::error(error), false);
}
let mut child = match spawn_with_etxtbsy_retry(|| {
self.build_command(run_provider_token, &policy, resume_id)
})
.await
{
Ok(child) => child,
Err(error) => {
return (
ChildOutcome::error(format!(
"spawn Codex CLI '{}': {error}; install with `npm i -g @openai/codex`, `brew install codex`, or an official GitHub release",
self.binary.display()
)),
false,
);
}
};
let Some(mut stdin) = child.stdin.take() else {
terminate_child(&mut child).await;
return (ChildOutcome::error("Codex child has no stdin pipe"), false);
};
let Some(stdout) = child.stdout.take() else {
terminate_child(&mut child).await;
return (ChildOutcome::error("Codex child has no stdout pipe"), false);
};
let stderr = child.stderr.take();
let write_result = tokio::select! {
result = async {
stdin.write_all(prompt.as_bytes()).await?;
stdin.write_all(b"\n").await?;
stdin.shutdown().await
} => result,
_ = cancel.cancelled() => {
terminate_child(&mut child).await;
events.emit(event_json(AgentEvent::Cancelled {
message: Some("Codex child cancelled".to_string()),
}));
return (ChildOutcome::cancelled(), false);
}
};
if let Err(error) = write_result {
terminate_child(&mut child).await;
return (
ChildOutcome::error(format!("write Codex prompt to stdin: {error}")),
false,
);
}
drop(stdin);
let stderr_tail = Arc::new(Mutex::new(String::new()));
let stderr_task = stderr.map(|stderr| {
let tail = stderr_tail.clone();
tokio::spawn(async move { drain_stderr_tail(stderr, tail).await })
});
let mut reader = BufReader::with_capacity(64 * 1024, stdout);
let mut state = RunState::default();
let mut read_error = None;
let mut cancelled = false;
loop {
tokio::select! {
_ = cancel.cancelled() => {
cancelled = true;
terminate_child(&mut child).await;
break;
}
line = read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES) => {
match line {
Ok(Some(bytes)) => {
if bytes.iter().all(u8::is_ascii_whitespace) {
continue;
}
match serde_json::from_slice::<Value>(&bytes) {
Ok(value) => {
if let Some(thread_id) = self.handle_event(
value,
&policy,
events,
&mut state,
) {
self.write_session_state(&thread_id).await;
}
}
Err(error) => tracing::debug!(%error, "codex: skipping unparsable stdout line"),
}
}
Ok(None) => break,
Err(error) => {
read_error = Some(format!("Codex stdout read error: {error}"));
terminate_child(&mut child).await;
break;
}
}
}
}
}
let status = if cancelled || read_error.is_some() {
child.try_wait().ok().flatten()
} else {
match tokio::time::timeout(PROCESS_EXIT_WAIT, child.wait()).await {
Ok(result) => result.ok(),
Err(_) => {
terminate_child(&mut child).await;
child.try_wait().ok().flatten()
}
}
};
if let Some(task) = stderr_task {
let _ = task.await;
}
let stderr = stderr_tail.lock().await.clone();
if cancelled {
events.emit(event_json(AgentEvent::Cancelled {
message: Some("Codex child cancelled".to_string()),
}));
return (ChildOutcome::cancelled(), false);
}
let exited_without_turn = !state.turn_started && !state.completed;
let outcome = if let Some(error) = read_error {
ChildOutcome::error(error)
} else if status.as_ref().is_some_and(|status| !status.success()) {
ChildOutcome::error(format!(
"Codex CLI exited with status {}; stderr tail: {}",
status
.as_ref()
.map(ToString::to_string)
.unwrap_or_else(|| "<unknown>".to_string()),
display_stderr_tail(&stderr)
))
} else if let Some(error) = state.failure {
ChildOutcome::error(format!(
"{error}; stderr tail: {}",
display_stderr_tail(&stderr)
))
} else if !state.completed {
ChildOutcome::error(format!(
"Codex CLI exited without a turn.completed event; stderr tail: {}",
display_stderr_tail(&stderr)
))
} else {
let final_text = if state.last_agent_message.trim().is_empty() {
self.read_last_message().await
} else {
Some(state.last_agent_message)
};
match final_text {
Some(text) => ChildOutcome::completed(text),
None => ChildOutcome::error(
"Codex CLI completed without a final agent message or output-last-message file",
),
}
};
(outcome, exited_without_turn)
}
fn emit_policy_bootstrap(&self, policy: &EffectiveCodexPolicy, events: &EventSink) {
events.emit(json!({
"type": "runner_progress",
"session_id": "codex",
"round_count": 0,
"executor": "codex",
"phase": "bootstrap",
"binary": self.binary,
"version": self.version,
"model": self.model,
"auth_mode": self.auth.mode().as_str(),
"codex_home_mode": self.codex_home_mode(),
"forward_env": self.forward_env,
"sandbox": policy.sandbox.as_str(),
"approval_policy": policy.approval_policy.as_str(),
"network_access": policy.network_access,
"policy_invocation": policy.invocation.as_str(),
"permission_profile": self.permissions.permission_profile.as_deref(),
}));
for warning in &policy.warnings {
tracing::warn!(message = %warning, "Codex spawn policy warning");
events.emit(json!({
"type": "runner_progress",
"session_id": "codex",
"round_count": 0,
"executor": "codex",
"phase": "policy_warning",
"level": "warning",
"message": warning,
"sandbox": policy.sandbox.as_str(),
"approval_policy": policy.approval_policy.as_str(),
"network_access": policy.network_access,
"policy_invocation": policy.invocation.as_str(),
}));
}
}
}
#[async_trait]
impl ChildExecutor for CodexExecutor {
async fn run(
&self,
spec: RunSpec,
events: EventSink,
mut steer: SteerInbox,
cancel: CancellationToken,
) -> ChildOutcome {
let steer_drain = tokio::spawn(async move { while steer.recv().await.is_some() {} });
let parent_bypass = spec
.permission_policy
.as_ref()
.map(|policy| policy.bypass_permissions)
.unwrap_or(self.permissions.provisioned_bypass);
if spec.messages.is_empty() {
self.delete_session_state().await;
}
let resume_id = if spec.messages.is_empty() {
None
} else {
self.resolve_resume_id().await
};
let body = if resume_id.is_some() || spec.messages.is_empty() {
spec.assignment.clone()
} else {
build_rehydrated_turn(&spec.messages, &spec.assignment)
};
let run_provider_token = spec
.secrets
.codex_provider_token
.as_ref()
.map(bamboo_subagent::proto::SecretValue::expose);
let (outcome, exited_without_turn) = self
.run_process(
&body,
resume_id.as_deref(),
run_provider_token,
parent_bypass,
&events,
&cancel,
)
.await;
let outcome = if resume_id.is_some() && exited_without_turn {
tracing::warn!(
"codex: resume exited before turn progress; retrying once with rehydrated history"
);
events.emit(json!({
"type": "runner_progress",
"session_id": "codex",
"round_count": 0,
"executor": "codex",
"phase": "resume_fallback",
"message": "resume failed before turn progress; retrying once with rehydrated history",
}));
self.delete_session_state().await;
let fallback_body = build_rehydrated_turn(&spec.messages, &spec.assignment);
self.run_process(
&fallback_body,
None,
run_provider_token,
parent_bypass,
&events,
&cancel,
)
.await
.0
} else {
outcome
};
steer_drain.abort();
outcome
}
}
#[derive(Default)]
struct RunState {
thread_id: String,
turn_started: bool,
completed: bool,
failure: Option<String>,
last_agent_message: String,
usage: TokenUsage,
tool_error_emitted: bool,
started_items: HashSet<String>,
item_text: HashMap<String, String>,
item_output: HashMap<String, String>,
}
impl RunState {
fn session_id(&self) -> String {
if self.thread_id.is_empty() {
"codex".to_string()
} else {
self.thread_id.clone()
}
}
}
fn handle_item(phase: &str, item: &Value, events: &EventSink, state: &mut RunState) {
let item_type = item.get("type").and_then(Value::as_str).unwrap_or("");
let item_id = item
.get("id")
.and_then(Value::as_str)
.filter(|id| !id.is_empty())
.map(str::to_string)
.unwrap_or_else(|| format!("codex-{item_type}"));
match item_type {
"agent_message" => {
let text = item.get("text").and_then(Value::as_str).unwrap_or("");
emit_text_delta(
&item_id,
text,
&mut state.item_text,
|delta| AgentEvent::Token { content: delta },
events,
);
if phase == "completed" && !text.is_empty() {
state.last_agent_message = text.to_string();
}
}
"reasoning" => {
let text = reasoning_text(item);
emit_text_delta(
&item_id,
&text,
&mut state.item_text,
|delta| AgentEvent::ReasoningToken { content: delta },
events,
);
}
"command_execution" => {
ensure_tool_started(
&item_id,
"Bash",
json!({ "command": item.get("command").cloned().unwrap_or(Value::Null) }),
events,
state,
);
let output = item
.get("aggregated_output")
.and_then(Value::as_str)
.unwrap_or("");
emit_tool_output_delta(&item_id, output, events, state);
if phase == "completed" {
let exit_code = item.get("exit_code").and_then(Value::as_i64);
let successful = exit_code == Some(0)
&& item.get("status").and_then(Value::as_str) != Some("failed");
if successful {
events.emit(event_json(AgentEvent::ToolComplete {
tool_call_id: item_id,
result: ToolResult::text(
true,
truncate_chars(output, TOOL_RESULT_TRUNCATE_CHARS),
),
}));
} else {
let error = if output.trim().is_empty() {
format!("command failed with exit code {exit_code:?}")
} else {
truncate_chars(output, TOOL_RESULT_TRUNCATE_CHARS)
};
events.emit(event_json(AgentEvent::ToolError {
tool_call_id: item_id,
error,
}));
state.tool_error_emitted = true;
}
}
}
"file_change" => {
let detail = item
.get("changes")
.or_else(|| item.get("patch"))
.cloned()
.unwrap_or_else(|| item.clone());
ensure_tool_started(
&item_id,
"ApplyPatch",
json!({ "changes": detail }),
events,
state,
);
if phase == "completed" {
complete_structured_tool(&item_id, item, events, state);
}
}
"mcp_tool_call" => {
let server = item.get("server").and_then(Value::as_str).unwrap_or("mcp");
let tool = item
.get("tool")
.or_else(|| item.get("name"))
.and_then(Value::as_str)
.unwrap_or("tool");
let tool_name = format!("{server}::{tool}");
let arguments = item
.get("arguments")
.or_else(|| item.get("input"))
.cloned()
.unwrap_or_else(|| json!({}));
ensure_tool_started(&item_id, &tool_name, arguments, events, state);
if phase == "completed" {
complete_structured_tool(&item_id, item, events, state);
}
}
"web_search" => {
let query = item.get("query").cloned().unwrap_or(Value::Null);
ensure_tool_started(
&item_id,
"WebSearch",
json!({ "query": query }),
events,
state,
);
if phase == "completed" {
complete_structured_tool(&item_id, item, events, state);
}
}
"todo_list" => {
events.emit(json!({
"type": "runner_progress",
"session_id": state.session_id(),
"round_count": 1,
"codex_item_type": "todo_list",
"item": item,
}));
}
other => {
tracing::debug!(item_type = other, phase, "codex: unrecognized item type");
}
}
}
fn ensure_tool_started(
item_id: &str,
tool_name: &str,
arguments: Value,
events: &EventSink,
state: &mut RunState,
) {
if state.started_items.insert(item_id.to_string()) {
events.emit(event_json(AgentEvent::ToolStart {
tool_call_id: item_id.to_string(),
tool_name: tool_name.to_string(),
arguments,
}));
}
}
fn emit_tool_output_delta(item_id: &str, output: &str, events: &EventSink, state: &mut RunState) {
let previous = state.item_output.entry(item_id.to_string()).or_default();
let delta = if output.starts_with(previous.as_str()) {
&output[previous.len()..]
} else {
output
};
if !delta.is_empty() {
events.emit(event_json(AgentEvent::ToolToken {
tool_call_id: item_id.to_string(),
content: delta.to_string(),
}));
}
*previous = output.to_string();
}
fn complete_structured_tool(item_id: &str, item: &Value, events: &EventSink, state: &mut RunState) {
if let Some(error) = item.get("error").filter(|value| !value.is_null()) {
events.emit(event_json(AgentEvent::ToolError {
tool_call_id: item_id.to_string(),
error: value_text(error),
}));
state.tool_error_emitted = true;
return;
}
let result = item
.get("result")
.or_else(|| item.get("output"))
.cloned()
.unwrap_or_else(|| item.clone());
events.emit(event_json(AgentEvent::ToolComplete {
tool_call_id: item_id.to_string(),
result: ToolResult::text(
true,
truncate_chars(&value_text(&result), TOOL_RESULT_TRUNCATE_CHARS),
),
}));
}
fn looks_like_sandbox_denial(text: &str) -> bool {
let normalized = text.to_ascii_lowercase();
let denied = normalized.contains("operation not permitted")
|| normalized.contains("permission denied")
|| normalized.contains("read-only file system");
let operation = normalized.contains("sandbox")
|| normalized.contains("command")
|| normalized.contains("write")
|| normalized.contains("exit status");
denied && operation
}
fn emit_text_delta<F>(
item_id: &str,
text: &str,
seen: &mut HashMap<String, String>,
build: F,
events: &EventSink,
) where
F: FnOnce(String) -> AgentEvent,
{
let previous = seen.entry(item_id.to_string()).or_default();
let delta = if text.starts_with(previous.as_str()) {
&text[previous.len()..]
} else {
text
};
if !delta.is_empty() {
events.emit(event_json(build(delta.to_string())));
}
*previous = text.to_string();
}
fn reasoning_text(item: &Value) -> String {
match item.get("text").or_else(|| item.get("summary")) {
Some(Value::String(text)) => text.clone(),
Some(Value::Array(parts)) => parts
.iter()
.filter_map(|part| {
part.as_str()
.or_else(|| part.get("text").and_then(Value::as_str))
})
.collect::<Vec<_>>()
.join("\n"),
Some(value) => value.to_string(),
None => String::new(),
}
}
fn parse_usage(value: Option<&Value>) -> TokenUsage {
let input = value
.and_then(|usage| usage.get("input_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
let output = value
.and_then(|usage| usage.get("output_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
TokenUsage {
prompt_tokens: input,
completion_tokens: output,
total_tokens: input.saturating_add(output),
}
}
fn final_text_from_terminal(value: &Value) -> Option<String> {
["final_output", "output_text", "result"]
.into_iter()
.find_map(|key| value.get(key).and_then(Value::as_str))
.filter(|text| !text.is_empty())
.map(str::to_string)
}
fn error_message(value: &Value, fallback: &str) -> String {
value
.pointer("/error/message")
.or_else(|| value.get("message"))
.or_else(|| value.get("error"))
.map(value_text)
.filter(|message| !message.is_empty())
.unwrap_or_else(|| fallback.to_string())
}
fn value_text(value: &Value) -> String {
value
.as_str()
.map(str::to_string)
.unwrap_or_else(|| value.to_string())
}
fn display_stderr_tail(stderr: &str) -> &str {
let trimmed = stderr.trim();
if trimmed.is_empty() {
"<empty>"
} else {
trimmed
}
}
fn event_json(event: AgentEvent) -> Value {
serde_json::to_value(event).unwrap_or_else(|_| json!({}))
}
fn truncate_chars(text: &str, max_chars: usize) -> String {
if text.chars().count() <= max_chars {
return text.to_string();
}
let head: String = text.chars().take(max_chars).collect();
let dropped = text.chars().count().saturating_sub(max_chars);
format!("{head}\n… [truncated, {dropped} more chars]")
}
fn has_git_metadata(path: &Path) -> bool {
path.ancestors()
.any(|ancestor| ancestor.join(".git").exists())
}
fn validate_codex_base_url(mode: CodexAuthMode, raw: Option<&str>) -> Result<(), String> {
match mode {
CodexAuthMode::Custom | CodexAuthMode::Bamboo => {
let raw = raw
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| format!("Codex auth mode '{mode:?}' requires codex_base_url"))?;
let parsed = url::Url::parse(raw)
.map_err(|error| format!("invalid Codex base URL '{raw}': {error}"))?;
if !matches!(parsed.scheme(), "http" | "https") || parsed.host_str().is_none() {
return Err("Codex base URL must be an absolute HTTP(S) URL".to_string());
}
if !parsed.username().is_empty()
|| parsed.password().is_some()
|| parsed.query().is_some()
|| parsed.fragment().is_some()
{
return Err(
"Codex base URL must not contain credentials, query parameters, or a fragment"
.to_string(),
);
}
}
CodexAuthMode::Inherit | CodexAuthMode::ApiKey => {
if raw.is_some_and(|value| !value.trim().is_empty()) {
return Err(
"codex_base_url is only valid with custom or bamboo auth mode".to_string(),
);
}
}
}
Ok(())
}
fn validate_codex_forward_env(mode: CodexAuthMode, names: &[String]) -> Result<(), String> {
let mut seen = HashSet::new();
for name in names {
let valid = name
.chars()
.next()
.is_some_and(|first| first == '_' || first.is_ascii_alphabetic())
&& name
.chars()
.all(|character| character == '_' || character.is_ascii_alphanumeric());
if !valid {
return Err(format!("invalid Codex forward_env name '{name}'"));
}
if !seen.insert(name.as_str()) {
return Err(format!("duplicate Codex forward_env name '{name}'"));
}
if name.starts_with("CODEX_") || name == CODEX_PROVIDER_ENV {
return Err(format!(
"Codex forward_env may not override reserved variable '{name}'"
));
}
}
let forwards_openai = seen.contains("OPENAI_API_KEY");
if mode == CodexAuthMode::ApiKey && !forwards_openai {
return Err(
"Codex api_key auth requires explicit OPENAI_API_KEY in codex_forward_env".to_string(),
);
}
if mode != CodexAuthMode::ApiKey && forwards_openai {
return Err("OPENAI_API_KEY may only be forwarded in Codex api_key auth mode".to_string());
}
Ok(())
}
async fn spawn_with_etxtbsy_retry(
mut build: impl FnMut() -> Result<Command, String>,
) -> Result<Child, String> {
let mut last_error = None;
for _ in 0..5 {
match build()?.spawn() {
Ok(child) => return Ok(child),
Err(error) if error.raw_os_error() == Some(26) => {
last_error = Some(error);
tokio::time::sleep(Duration::from_millis(10)).await;
}
Err(error) => return Err(error.to_string()),
}
}
Err(last_error
.expect("retry loop records ETXTBSY before exhausting")
.to_string())
}
#[derive(Clone, Copy)]
enum ProcessSignal {
Term,
Kill,
}
#[cfg(unix)]
fn signal_process_group(pgid: libc::pid_t, signal: ProcessSignal) {
let signal = match signal {
ProcessSignal::Term => libc::SIGTERM,
ProcessSignal::Kill => libc::SIGKILL,
};
unsafe {
libc::kill(-pgid, signal);
}
}
#[cfg(unix)]
fn signal_process(pid: libc::pid_t, signal: ProcessSignal) {
let signal = match signal {
ProcessSignal::Term => libc::SIGTERM,
ProcessSignal::Kill => libc::SIGKILL,
};
unsafe {
libc::kill(pid, signal);
}
}
#[cfg(unix)]
fn descendant_processes(root: libc::pid_t) -> Vec<libc::pid_t> {
use sysinfo::{ProcessesToUpdate, System};
let mut system = System::new();
system.refresh_processes(ProcessesToUpdate::All, true);
let mut known = HashSet::from([root]);
let mut descendants = Vec::new();
loop {
let mut added = false;
for (pid, process) in system.processes() {
let pid = pid.as_u32() as libc::pid_t;
let parent = process
.parent()
.map(|parent| parent.as_u32() as libc::pid_t);
if !known.contains(&pid) && parent.is_some_and(|parent| known.contains(&parent)) {
known.insert(pid);
descendants.push(pid);
added = true;
}
}
if !added {
break;
}
}
descendants
}
#[cfg(unix)]
fn process_exists(pid: libc::pid_t) -> bool {
if unsafe { libc::kill(pid, 0) } == 0 {
return true;
}
std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
}
#[cfg(unix)]
fn process_group_exists(pgid: libc::pid_t) -> bool {
if unsafe { libc::kill(-pgid, 0) } == 0 {
return true;
}
std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
}
#[cfg(unix)]
pub(crate) async fn terminate_child(child: &mut Child) {
let Some(pgid) = child.id().map(|pid| pid as libc::pid_t) else {
return;
};
let descendants = descendant_processes(pgid);
for &pid in descendants.iter().rev() {
signal_process(pid, ProcessSignal::Term);
}
signal_process_group(pgid, ProcessSignal::Term);
let deadline = tokio::time::Instant::now() + SIGTERM_WAIT;
loop {
let _ = child.try_wait();
if !process_group_exists(pgid) && !descendants.iter().copied().any(process_exists) {
let _ = child.wait().await;
return;
}
if tokio::time::Instant::now() >= deadline {
break;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
for &pid in descendants.iter().rev() {
signal_process(pid, ProcessSignal::Kill);
}
signal_process_group(pgid, ProcessSignal::Kill);
let _ = child.start_kill();
let _ = child.wait().await;
}
#[cfg(not(unix))]
pub(crate) async fn terminate_child(child: &mut Child) {
if tokio::time::timeout(SIGTERM_WAIT, child.wait())
.await
.is_ok()
{
return;
}
let _ = child.start_kill();
let _ = child.wait().await;
}
pub(crate) async fn read_bounded_line<R>(
reader: &mut R,
max_bytes: usize,
) -> std::io::Result<Option<Vec<u8>>>
where
R: tokio::io::AsyncBufRead + Unpin,
{
let mut output = Vec::new();
loop {
let (found_newline, consumed) = {
let available = reader.fill_buf().await?;
if available.is_empty() {
return Ok(if output.is_empty() {
None
} else {
Some(output)
});
}
match available.iter().position(|byte| *byte == b'\n') {
Some(position) => {
output.extend_from_slice(&available[..position]);
(true, position + 1)
}
None => {
output.extend_from_slice(available);
(false, available.len())
}
}
};
reader.consume(consumed);
if output.len() > max_bytes {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("stdout line exceeded {max_bytes} bytes"),
));
}
if found_newline {
return Ok(Some(output));
}
}
}
async fn drain_stderr_tail(stderr: tokio::process::ChildStderr, tail: Arc<Mutex<String>>) {
let mut reader = BufReader::new(stderr);
let mut buffer = Vec::new();
loop {
buffer.clear();
match reader.read_until(b'\n', &mut buffer).await {
Ok(0) | Err(_) => return,
Ok(_) => {
let mut tail = tail.lock().await;
tail.push_str(&String::from_utf8_lossy(&buffer));
if tail.len() > STDERR_TAIL_BYTES {
let excess = tail.len() - STDERR_TAIL_BYTES;
let cut = tail
.char_indices()
.map(|(index, _)| index)
.find(|index| *index >= excess)
.unwrap_or(tail.len());
tail.drain(..cut);
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use bamboo_subagent::proto::TerminalStatus;
use bamboo_subagent::provision::ScopedCredential;
fn resolved_auth(
mode: Option<&str>,
base_url: Option<&str>,
provider_key_ref: Option<&str>,
credentials: &[ScopedCredential],
forward_env: &[&str],
) -> Result<CodexAuthConfig, String> {
resolve_codex_auth_config(
mode,
false,
base_url.map(str::to_string),
Some("responses".to_string()),
provider_key_ref,
credentials,
&forward_env
.iter()
.map(|name| name.to_string())
.collect::<Vec<_>>(),
)
}
fn command_env(command: &Command, name: &str) -> Option<String> {
command
.as_std()
.get_envs()
.find(|(key, _)| *key == name)
.and_then(|(_, value)| value)
.map(|value| value.to_string_lossy().into_owned())
}
fn command_args(command: &Command) -> Vec<String> {
command
.as_std()
.get_args()
.map(|argument| argument.to_string_lossy().into_owned())
.collect()
}
fn auth_error(result: Result<CodexAuthConfig, String>) -> String {
match result {
Err(error) => error,
Ok(_) => panic!("expected Codex auth configuration to be rejected"),
}
}
fn permissions(
sandbox: Option<&str>,
approval_policy: Option<&str>,
network_access: bool,
allow_danger_bypass: bool,
permission_profile: Option<&str>,
provisioned_bypass: bool,
workspace_owned: bool,
) -> CodexPermissionConfig {
resolve_codex_permission_config(
sandbox,
approval_policy,
network_access,
allow_danger_bypass,
permission_profile.map(str::to_string),
provisioned_bypass,
workspace_owned,
)
.unwrap()
}
fn default_permissions() -> CodexPermissionConfig {
permissions(None, None, false, false, None, false, false)
}
fn default_policy(executor: &CodexExecutor) -> EffectiveCodexPolicy {
executor.permissions.effective(false, false)
}
fn fixture_executor() -> CodexExecutor {
CodexExecutor {
binary: PathBuf::from("/usr/local/bin/codex"),
version: "codex-cli 0.144.5".to_string(),
model: None,
permissions: default_permissions(),
workspace: None,
state_dir: None,
forward_env: Vec::new(),
auth: CodexAuthConfig::inherit(),
}
}
fn map_fixture(input: &str) -> (RunState, Vec<Value>) {
let executor = fixture_executor();
let (sink, mut rx) = EventSink::channel();
let mut state = RunState::default();
let policy = default_policy(&executor);
for line in input.lines().filter(|line| !line.trim().is_empty()) {
executor.handle_event(
serde_json::from_str(line).unwrap(),
&policy,
&sink,
&mut state,
);
}
let mut events = Vec::new();
while let Ok(event) = rx.try_recv() {
events.push(event);
}
(state, events)
}
#[test]
fn permission_profile_mapping_table_is_explicit_and_double_gated() {
struct Case {
name: &'static str,
permissions: CodexPermissionConfig,
parent_bypass: bool,
is_root: bool,
sandbox: CodexSandbox,
approval: CodexApprovalPolicy,
invocation: CodexPolicyInvocation,
network: bool,
warning: bool,
}
let cases = [
Case {
name: "default",
permissions: permissions(None, None, false, false, None, false, false),
parent_bypass: false,
is_root: false,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::Explicit,
network: false,
warning: false,
},
Case {
name: "restricted",
permissions: permissions(
None,
None,
false,
false,
Some("restricted"),
false,
false,
),
parent_bypass: false,
is_root: false,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::Explicit,
network: false,
warning: false,
},
Case {
name: "research",
permissions: permissions(None, None, false, false, Some("research"), false, false),
parent_bypass: false,
is_root: false,
sandbox: CodexSandbox::ReadOnly,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::Explicit,
network: false,
warning: false,
},
Case {
name: "network workspace",
permissions: permissions(None, None, true, false, None, false, false),
parent_bypass: false,
is_root: false,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::Explicit,
network: true,
warning: false,
},
Case {
name: "ordinary bypass stays sandboxed",
permissions: permissions(None, None, false, false, None, false, false),
parent_bypass: true,
is_root: false,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::FullAuto,
network: false,
warning: false,
},
Case {
name: "config-only danger opt-in downgrades",
permissions: permissions(None, None, false, true, None, false, false),
parent_bypass: false,
is_root: false,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::FullAuto,
network: false,
warning: true,
},
Case {
name: "both danger gates",
permissions: permissions(None, None, false, true, None, false, false),
parent_bypass: true,
is_root: false,
sandbox: CodexSandbox::DangerFullAccess,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::DangerBypass,
network: true,
warning: true,
},
Case {
name: "root danger request downgrades",
permissions: permissions(
Some("danger-full-access"),
None,
false,
true,
None,
false,
false,
),
parent_bypass: true,
is_root: true,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::Never,
invocation: CodexPolicyInvocation::FullAuto,
network: false,
warning: true,
},
Case {
name: "explicit safe override",
permissions: permissions(
Some("workspace-write"),
Some("on-failure"),
false,
false,
Some("read-only"),
false,
false,
),
parent_bypass: false,
is_root: false,
sandbox: CodexSandbox::WorkspaceWrite,
approval: CodexApprovalPolicy::OnFailure,
invocation: CodexPolicyInvocation::Explicit,
network: false,
warning: false,
},
];
for case in cases {
let actual = case.permissions.effective(case.parent_bypass, case.is_root);
assert_eq!(actual.sandbox, case.sandbox, "{} sandbox", case.name);
assert_eq!(
actual.approval_policy, case.approval,
"{} approval",
case.name
);
assert_eq!(
actual.invocation, case.invocation,
"{} invocation",
case.name
);
assert_eq!(actual.network_access, case.network, "{} network", case.name);
assert_eq!(
!actual.warnings.is_empty(),
case.warning,
"{} warning",
case.name
);
}
}
#[test]
fn permission_configuration_rejects_interactive_and_incoherent_values() {
assert!(resolve_codex_permission_config(
None,
Some("on-request"),
false,
false,
None,
false,
false,
)
.unwrap_err()
.contains("non-interactive"));
assert!(resolve_codex_permission_config(
Some("read-only"),
None,
true,
false,
None,
false,
false,
)
.unwrap_err()
.contains("effective workspace-write"));
assert!(resolve_codex_permission_config(
None,
None,
true,
false,
Some("research".to_string()),
false,
false,
)
.unwrap_err()
.contains("effective workspace-write"));
}
#[test]
fn command_flags_match_effective_policy_and_workspace_ownership() {
let workspace = tempfile::tempdir().unwrap();
let mut executor = fixture_executor();
executor.workspace = Some(workspace.path().to_string_lossy().into_owned());
let policy = executor.permissions.effective(false, false);
let args = command_args(&executor.build_command(None, &policy, None).unwrap());
assert!(args
.windows(2)
.any(|pair| pair == ["--sandbox", "workspace-write"]));
assert!(args
.windows(2)
.any(|pair| pair == ["--config", "approval_policy=\"never\""]));
assert!(!args
.iter()
.any(|argument| argument == "--skip-git-repo-check"));
executor.permissions = permissions(None, None, false, false, None, false, true);
let policy = executor.permissions.effective(false, false);
let args = command_args(&executor.build_command(None, &policy, None).unwrap());
assert!(args
.iter()
.any(|argument| argument == "--skip-git-repo-check"));
executor.permissions = permissions(None, None, true, false, None, false, true);
let policy = executor.permissions.effective(false, false);
let args = command_args(&executor.build_command(None, &policy, None).unwrap());
assert!(args
.windows(2)
.any(|pair| { pair == ["--config", "sandbox_workspace_write.network_access=true"] }));
executor.permissions = permissions(None, None, false, false, None, false, true);
let policy = executor.permissions.effective(true, false);
let args = command_args(&executor.build_command(None, &policy, None).unwrap());
assert!(args.iter().any(|argument| argument == "--full-auto"));
assert!(!args
.iter()
.any(|argument| argument == "--dangerously-bypass-approvals-and-sandbox"));
executor.permissions = permissions(
Some("danger-full-access"),
None,
false,
true,
None,
false,
true,
);
let policy = executor.permissions.effective(true, false);
let args = command_args(&executor.build_command(None, &policy, None).unwrap());
assert!(args
.iter()
.any(|argument| argument == "--dangerously-bypass-approvals-and-sandbox"));
assert!(!args.iter().any(|argument| argument == "--full-auto"));
}
#[test]
fn danger_bypass_emits_loud_audit_warning() {
let mut executor = fixture_executor();
executor.permissions = permissions(
Some("danger-full-access"),
None,
false,
true,
Some("bypass"),
false,
false,
);
let policy = executor.permissions.effective(true, false);
let (sink, mut receiver) = EventSink::channel();
executor.emit_policy_bootstrap(&policy, &sink);
let events = std::iter::from_fn(|| receiver.try_recv().ok()).collect::<Vec<_>>();
assert!(events.iter().any(|event| {
event["phase"] == "bootstrap"
&& event["sandbox"] == "danger-full-access"
&& event["approval_policy"] == "never"
&& event["policy_invocation"] == "danger-bypass"
}));
assert!(events.iter().any(|event| {
event["phase"] == "policy_warning"
&& event["level"] == "warning"
&& event["message"]
.as_str()
.is_some_and(|message| message.contains("DANGER"))
}));
}
#[test]
fn auth_mode_matrix_is_explicit_isolated_and_secret_safe() {
let credential = ScopedCredential {
provider: "custom-provider".to_string(),
api_key: "custom-secret-570".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.custom.api_key".to_string()),
};
let inherit = resolved_auth(Some("inherit"), None, None, &[], &[]).unwrap();
assert_eq!(inherit.mode(), CodexAuthMode::Inherit);
let mut inherit_executor = fixture_executor();
inherit_executor.auth = inherit;
let policy = default_policy(&inherit_executor);
let inherit_command = inherit_executor.build_command(None, &policy, None).unwrap();
assert!(command_env(&inherit_command, "CODEX_HOME").is_none());
let inherit_args = inherit_command
.as_std()
.get_args()
.map(|arg| arg.to_string_lossy().into_owned())
.collect::<Vec<_>>();
assert!(!inherit_args.iter().any(|arg| arg == "--ignore-rules"));
let api_key = resolved_auth(Some("api_key"), None, None, &[], &["OPENAI_API_KEY"]).unwrap();
assert_eq!(api_key.mode(), CodexAuthMode::ApiKey);
assert!(
auth_error(resolved_auth(Some("api_key"), None, None, &[], &[]))
.contains("explicit OPENAI_API_KEY")
);
let custom = resolved_auth(
Some("custom"),
Some("https://provider.example/v1"),
Some("provider.custom.api_key"),
std::slice::from_ref(&credential),
&[],
)
.unwrap();
assert_eq!(custom.mode(), CodexAuthMode::Custom);
let generated = custom.generated_config_toml().unwrap();
assert!(generated.contains("model_provider = \"custom\""));
assert!(generated.contains("base_url = \"https://provider.example/v1\""));
assert!(generated.contains("env_key = \"BAMBOO_CODEX_PROVIDER_KEY\""));
assert!(generated.contains("wire_api = \"responses\""));
assert!(!generated.contains("custom-secret-570"));
let state = tempfile::tempdir().unwrap();
let mut custom_executor = fixture_executor();
custom_executor.state_dir = Some(state.path().to_path_buf());
custom_executor.auth = custom;
let policy = default_policy(&custom_executor);
let custom_command = custom_executor.build_command(None, &policy, None).unwrap();
assert_eq!(
command_env(&custom_command, CODEX_PROVIDER_ENV).as_deref(),
Some("custom-secret-570")
);
assert_eq!(
command_env(&custom_command, "CODEX_HOME"),
Some(
state
.path()
.join("codex-home")
.to_string_lossy()
.into_owned()
)
);
assert!(custom_command
.as_std()
.get_args()
.any(|arg| arg == "--ignore-rules"));
let bamboo = resolved_auth(
None,
Some("http://127.0.0.1:9562/openai/v1"),
None,
&[],
&[],
)
.unwrap();
assert_eq!(bamboo.mode(), CodexAuthMode::Bamboo);
let generated = bamboo.generated_config_toml().unwrap();
assert!(generated.contains("model_provider = \"bamboo\""));
assert!(generated.contains("http://127.0.0.1:9562/openai/v1"));
assert!(!generated.contains("bcx1_"));
let app_server_generated = bamboo
.generated_app_server_config_toml(
Path::new("/opt/bamboo/bin/bamboo"),
Path::new("/private/state/codex-provider-token"),
)
.unwrap();
assert!(app_server_generated.contains("[model_providers.bamboo.auth]"));
assert!(app_server_generated.contains("command = \"/opt/bamboo/bin/bamboo\""));
assert!(app_server_generated.contains("\"codex-provider-token\""));
assert!(app_server_generated.contains("refresh_interval_ms = 1"));
assert!(!app_server_generated.contains("env_key"));
assert!(!app_server_generated.contains("bcx1_"));
let mut bamboo_executor = fixture_executor();
bamboo_executor.state_dir = Some(state.path().to_path_buf());
bamboo_executor.auth = bamboo;
let policy = default_policy(&bamboo_executor);
assert!(bamboo_executor
.build_command(None, &policy, None)
.unwrap_err()
.contains("per-run provider token"));
let bamboo_command = bamboo_executor
.build_command(Some("bcx1_per_run_570"), &policy, None)
.unwrap();
assert_eq!(
command_env(&bamboo_command, CODEX_PROVIDER_ENV).as_deref(),
Some("bcx1_per_run_570")
);
}
#[test]
fn auth_matrix_rejects_ambiguous_or_unsafe_configuration() {
for mode in ["inherit", "custom", "bamboo"] {
let error = auth_error(resolved_auth(
Some(mode),
match mode {
"custom" => Some("https://provider.example/v1"),
"bamboo" => Some("http://127.0.0.1:9562/openai/v1"),
_ => None,
},
(mode == "custom").then_some("provider.missing.api_key"),
&[],
&["OPENAI_API_KEY"],
));
assert!(
error.contains("only be forwarded") || error.contains("did not resolve"),
"mode {mode}: {error}"
);
}
assert!(
auth_error(resolved_auth(Some("future"), None, None, &[], &[]))
.contains("unknown Codex auth mode")
);
assert!(auth_error(resolve_codex_auth_config(
Some("inherit"),
false,
None,
Some("chat".to_string()),
None,
&[],
&[],
))
.contains("requires responses"));
assert!(auth_error(resolved_auth(
Some("inherit"),
None,
None,
&[],
&["CODEX_HOME"],
))
.contains("reserved variable"));
assert!(auth_error(resolved_auth(
Some("custom"),
Some("https://user:secret@provider.example/v1?x=1"),
Some("provider.custom.api_key"),
&[],
&[],
))
.contains("must not contain credentials"));
}
#[cfg(unix)]
#[tokio::test]
async fn isolated_home_removes_stale_login_and_writes_locked_secret_free_config() {
use std::os::unix::fs::PermissionsExt as _;
let state = tempfile::tempdir().unwrap();
let home = state.path().join("codex-home");
std::fs::create_dir_all(&home).unwrap();
std::fs::write(home.join("auth.json"), r#"{"tokens":"stale-secret"}"#).unwrap();
let mut executor = fixture_executor();
executor.state_dir = Some(state.path().to_path_buf());
executor.auth = CodexAuthConfig {
mode: CodexAuthMode::Custom,
base_url: Some("https://provider.example/v1".to_string()),
wire_api: "responses".to_string(),
provider_key: Some("custom-secret-570".to_string()),
};
executor.prepare_auth_home().await.unwrap();
assert!(!home.join("auth.json").exists());
let config_path = home.join("config.toml");
let config = std::fs::read_to_string(&config_path).unwrap();
assert!(config.contains("https://provider.example/v1"));
assert!(!config.contains("custom-secret-570"));
assert_eq!(
std::fs::metadata(&home).unwrap().permissions().mode() & 0o777,
0o700
);
assert_eq!(
std::fs::metadata(&config_path)
.unwrap()
.permissions()
.mode()
& 0o777,
0o600
);
std::fs::write(home.join("auth.json"), r#"{"tokens":"late-secret"}"#).unwrap();
std::fs::write(&config_path, "model_provider = \"attacker\"\n").unwrap();
executor.prepare_auth_home().await.unwrap();
assert!(!home.join("auth.json").exists());
let restored = std::fs::read_to_string(&config_path).unwrap();
assert!(restored.contains("model_provider = \"custom\""));
assert!(!restored.contains("attacker"));
}
#[tokio::test]
async fn bounded_reader_accepts_ten_megabytes_and_rejects_more() {
let allowed = vec![b'x'; MAX_STDOUT_LINE_BYTES];
let mut bytes = allowed.clone();
bytes.push(b'\n');
let mut reader = BufReader::new(bytes.as_slice());
assert_eq!(
read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES)
.await
.unwrap()
.unwrap()
.len(),
MAX_STDOUT_LINE_BYTES
);
let mut too_large = vec![b'x'; MAX_STDOUT_LINE_BYTES + 1];
too_large.push(b'\n');
let mut reader = BufReader::new(too_large.as_slice());
assert_eq!(
read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES)
.await
.unwrap_err()
.kind(),
std::io::ErrorKind::InvalidData
);
}
#[test]
fn usage_matches_real_codex_schema() {
let usage = parse_usage(Some(&json!({
"input_tokens": 13_460,
"cached_input_tokens": 9_984,
"output_tokens": 6,
"reasoning_output_tokens": 0
})));
assert_eq!(usage.prompt_tokens, 13_460);
assert_eq!(usage.completion_tokens, 6);
assert_eq!(usage.total_tokens, 13_466);
}
#[test]
fn recorded_jsonl_fixtures_map_agent_and_command_events() {
let cases = [
(
include_str!("../tests/fixtures/codex-cli/0.144.5-simple.jsonl"),
"PONG",
false,
13_466,
),
(
include_str!("../tests/fixtures/codex-cli/0.144.5-command.jsonl"),
"The current working directory is `/private/tmp/zenith-bamboo-569-codex-cli`.",
true,
27_108,
),
];
for (fixture, expected_final, expects_tool, total_tokens) in cases {
let (state, events) = map_fixture(fixture);
assert!(state.completed);
assert_eq!(state.last_agent_message, expected_final);
assert_eq!(state.usage.total_tokens, total_tokens);
assert!(events.iter().any(|event| event["type"] == "token"));
assert!(events.iter().any(|event| event["type"] == "complete"));
assert_eq!(
events.iter().any(|event| event["type"] == "tool_start"),
expects_tool
);
if expects_tool {
assert!(events.iter().any(|event| {
event["type"] == "tool_start" && event["tool_name"] == "Bash"
}));
assert!(events.iter().any(|event| event["type"] == "tool_token"));
assert!(events.iter().any(|event| event["type"] == "tool_complete"));
}
}
}
#[test]
fn item_mapping_table_covers_non_command_codex_items() {
let cases = [
(
json!({"id":"r1","type":"reasoning","text":"thinking"}),
"reasoning_token",
None,
),
(
json!({"id":"f1","type":"file_change","changes":[{"path":"a.txt","kind":"add"}]}),
"tool_start",
Some("ApplyPatch"),
),
(
json!({"id":"m1","type":"mcp_tool_call","server":"files","tool":"read","arguments":{"path":"a.txt"},"result":"ok"}),
"tool_start",
Some("files::read"),
),
(
json!({"id":"w1","type":"web_search","query":"Bamboo","result":[]}),
"tool_start",
Some("WebSearch"),
),
(
json!({"id":"t1","type":"todo_list","items":[]}),
"runner_progress",
None,
),
];
for (item, expected_type, expected_tool) in cases {
let (sink, mut rx) = EventSink::channel();
let mut state = RunState::default();
handle_item("completed", &item, &sink, &mut state);
let first = rx.try_recv().expect("mapping emitted an event");
assert_eq!(first["type"], expected_type, "item: {item}");
if let Some(tool) = expected_tool {
assert_eq!(first["tool_name"], tool, "item: {item}");
assert!(rx
.try_recv()
.is_ok_and(|event| event["type"] == "tool_complete"));
}
}
}
#[test]
fn unknown_event_and_item_types_are_tolerated() {
let executor = fixture_executor();
let (sink, mut rx) = EventSink::channel();
let mut state = RunState::default();
let policy = default_policy(&executor);
executor.handle_event(
json!({"type":"future.event","payload":{"schema":2}}),
&policy,
&sink,
&mut state,
);
executor.handle_event(
json!({"type":"item.completed","item":{"id":"x","type":"future_item"}}),
&policy,
&sink,
&mut state,
);
assert!(rx.try_recv().is_err());
assert!(!state.completed);
assert!(state.failure.is_none());
}
#[test]
fn sandbox_denial_report_without_cli_tool_item_becomes_tool_error() {
let executor = fixture_executor();
let policy = executor.permissions.effective(false, false);
let (sink, mut receiver) = EventSink::channel();
let mut state = RunState::default();
executor.handle_event(
json!({
"type": "item.completed",
"item": {
"id": "denied-report",
"type": "agent_message",
"text": "Command exited with status 1. Sandbox failure: outside write: Operation not permitted"
}
}),
&policy,
&sink,
&mut state,
);
let events = std::iter::from_fn(|| receiver.try_recv().ok()).collect::<Vec<_>>();
assert!(events.iter().any(|event| event["type"] == "token"));
assert!(events.iter().any(|event| {
event["type"] == "tool_error"
&& event["tool_call_id"] == "denied-report-sandbox-denial"
&& event["error"]
.as_str()
.is_some_and(|error| error.contains("Operation not permitted"))
}));
}
#[cfg(unix)]
#[test]
fn provider_token_reader_checks_open_descriptor_and_rejects_symlinks() {
use std::os::unix::fs::{symlink, PermissionsExt as _};
let root = tempfile::tempdir().unwrap();
let token = root.path().join("token");
std::fs::write(&token, " bcx1_short_lived \n").unwrap();
std::fs::set_permissions(&token, std::fs::Permissions::from_mode(0o600)).unwrap();
assert_eq!(
read_codex_provider_token(&token).unwrap(),
"bcx1_short_lived"
);
let link = root.path().join("token-link");
symlink(&token, &link).unwrap();
assert!(read_codex_provider_token(&link).is_err());
std::fs::set_permissions(&token, std::fs::Permissions::from_mode(0o640)).unwrap();
assert!(read_codex_provider_token(&token)
.unwrap_err()
.contains("group/other"));
}
#[cfg(unix)]
mod unix_process_tests {
use super::*;
use std::io::Write as _;
use std::os::unix::fs::PermissionsExt;
fn write_stub(dir: &Path, body: &str) -> PathBuf {
let path = dir.join("codex");
let mut file = std::fs::File::create(&path).unwrap();
writeln!(file, "#!/bin/sh").unwrap();
file.write_all(body.as_bytes()).unwrap();
let mut permissions = std::fs::metadata(&path).unwrap().permissions();
permissions.set_mode(0o755);
std::fs::set_permissions(&path, permissions).unwrap();
path
}
fn executor(binary: PathBuf, workspace: &Path) -> CodexExecutor {
CodexExecutor {
binary,
version: "codex-cli 0.144.5".to_string(),
model: None,
permissions: permissions(None, None, false, false, None, false, true),
workspace: Some(workspace.to_string_lossy().into_owned()),
state_dir: None,
forward_env: Vec::new(),
auth: CodexAuthConfig::inherit(),
}
}
fn run_spec(assignment: &str) -> RunSpec {
RunSpec {
assignment: assignment.to_string(),
logical_session: None,
project_id: None,
reasoning_effort: None,
permission_policy: None,
messages: Vec::new(),
activation_run_id: None,
initial_session_messages: Vec::new(),
secrets: Default::default(),
}
}
fn run_spec_with_messages(assignment: &str, messages: Vec<Value>) -> RunSpec {
RunSpec {
assignment: assignment.to_string(),
logical_session: None,
project_id: None,
reasoning_effort: None,
permission_policy: None,
messages,
activation_run_id: None,
initial_session_messages: Vec::new(),
secrets: Default::default(),
}
}
fn message(role: &str, content: &str) -> Value {
json!({"role": role, "content": content})
}
#[tokio::test]
async fn preflight_accepts_current_surface_and_rejects_old_or_missing_binary() {
let dir = tempfile::tempdir().unwrap();
let current = write_stub(
dir.path(),
r#"
case "$*" in
"--version") echo 'codex-cli 0.144.5' ;;
"exec --help") echo '--json --output-last-message --config --sandbox --dangerously-bypass-approvals-and-sandbox prompt from stdin' ;;
"exec resume --help") echo '--json' ;;
*) exit 2 ;;
esac
"#,
);
let checked = CodexExecutor::new(
Some(current.to_string_lossy().into_owned()),
None,
None,
None,
Vec::new(),
CodexAuthConfig::inherit(),
default_permissions(),
)
.await
.unwrap();
assert_eq!(checked.version, "codex-cli 0.144.5");
let old_dir = tempfile::tempdir().unwrap();
let old = write_stub(
old_dir.path(),
r#"
if [ "$1" = "--version" ]; then echo 'codex-cli 0.143.9'; else exit 2; fi
"#,
);
let old_error = CodexExecutor::new(
Some(old.to_string_lossy().into_owned()),
None,
None,
None,
Vec::new(),
CodexAuthConfig::inherit(),
default_permissions(),
)
.await
.err()
.expect("old version rejected");
assert!(old_error.contains("too old"));
assert!(old_error.contains(">= 0.144.0"));
let missing_error = CodexExecutor::new(
Some("/definitely/missing/codex".to_string()),
None,
None,
None,
Vec::new(),
CodexAuthConfig::inherit(),
default_permissions(),
)
.await
.err()
.expect("missing binary rejected");
assert!(missing_error.contains("npm i -g @openai/codex"));
assert!(missing_error.contains("brew install codex"));
assert!(missing_error.contains("codex_binary"));
}
#[tokio::test]
async fn clean_completion_and_nonzero_exit_have_distinct_outcomes() {
let workspace = tempfile::tempdir().unwrap();
let ok_dir = tempfile::tempdir().unwrap();
let ok = write_stub(
ok_dir.path(),
r#"
read -r prompt
echo '{"type":"thread.started","thread_id":"stub-thread"}'
echo '{"type":"turn.started"}'
echo '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"PONG"}}'
echo '{"type":"turn.completed","usage":{"input_tokens":3,"output_tokens":1}}'
"#,
);
let (sink, _rx) = EventSink::channel();
let outcome = executor(ok, workspace.path())
.run(
run_spec("reply PONG"),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Completed);
assert_eq!(outcome.result.as_deref(), Some("PONG"));
let fail_dir = tempfile::tempdir().unwrap();
let fail = write_stub(
fail_dir.path(),
r#"
read -r prompt
echo 'credential lookup failed' >&2
exit 7
"#,
);
let (sink, _rx) = EventSink::channel();
let outcome = executor(fail, workspace.path())
.run(
run_spec("fail"),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Error);
let error = outcome.error.unwrap();
assert!(error.contains("status exit status: 7"), "{error}");
assert!(error.contains("credential lookup failed"), "{error}");
}
#[tokio::test]
async fn spawn_uses_stdin_safe_defaults_and_last_message_fallback() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let state_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
: > "$DIR/argv.txt"
OUT=''
while [ "$#" -gt 0 ]; do
printf '%s\n' "$1" >> "$DIR/argv.txt"
if [ "$1" = '--output-last-message' ]; then
shift
OUT="$1"
printf '%s\n' "$1" >> "$DIR/argv.txt"
fi
shift
done
IFS= read -r prompt
printf '%s\n' "$prompt" > "$DIR/stdin.txt"
printf 'PONG\n' > "$OUT"
echo '{"type":"thread.started","thread_id":"fallback-thread"}'
echo '{"type":"turn.started"}'
echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
"#,
);
let mut exec = executor(bin, workspace.path());
exec.state_dir = Some(state_dir.path().to_path_buf());
let assignment = "a private prompt that must not appear in argv";
let (sink, _rx) = EventSink::channel();
let outcome = exec
.run(
run_spec(assignment),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Completed);
assert_eq!(outcome.result.as_deref(), Some("PONG"));
assert_eq!(
std::fs::read_to_string(bin_dir.path().join("stdin.txt"))
.unwrap()
.trim(),
assignment
);
let argv = std::fs::read_to_string(bin_dir.path().join("argv.txt")).unwrap();
assert!(
!argv.contains(assignment),
"prompt leaked into argv: {argv}"
);
for required in [
"exec",
"--json",
"--cd",
"--skip-git-repo-check",
"--sandbox",
"workspace-write",
"--config",
"approval_policy=\"never\"",
"--output-last-message",
"-",
] {
assert!(
argv.lines().any(|arg| arg == required),
"missing {required}: {argv}"
);
}
assert!(!argv.lines().any(|arg| arg == "--ignore-rules"));
assert!(!argv.lines().any(|arg| arg == "--ignore-user-config"));
}
#[tokio::test]
async fn thread_state_is_captured_resumed_recaptured_and_cleared_for_fresh_run() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let state_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
N=$(cat "$DIR/count" 2>/dev/null || echo 0)
N=$((N+1))
echo "$N" > "$DIR/count"
printf '%s\n' "$@" > "$DIR/argv-$N.txt"
IFS= read -r prompt
printf '%s\n' "$prompt" > "$DIR/stdin-$N.txt"
case "$N" in
1) THREAD='thread-one'; ANSWER='first' ;;
2) THREAD='thread-two'; ANSWER='second' ;;
*) THREAD=''; ANSWER='third' ;;
esac
if [ -n "$THREAD" ]; then
printf '{"type":"thread.started","thread_id":"%s"}\n' "$THREAD"
fi
echo '{"type":"turn.started"}'
printf '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"%s"}}\n' "$ANSWER"
echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
"#,
);
let mut exec = executor(bin, workspace.path());
exec.state_dir = Some(state_dir.path().to_path_buf());
let state_path = state_dir.path().join(CODEX_SESSION_STATE_FILE);
let (sink, _rx) = EventSink::channel();
let first = exec
.run(
run_spec("remember amber-572"),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(first.status, TerminalStatus::Completed);
assert!(!std::fs::read_to_string(bin_dir.path().join("argv-1.txt"))
.unwrap()
.lines()
.any(|argument| argument == "resume"));
let state: CodexSessionState =
serde_json::from_slice(&std::fs::read(&state_path).unwrap()).unwrap();
assert_eq!(state.thread_id, "thread-one");
assert_eq!(state.workspace, exec.workspace);
assert_eq!(state.codex_home_mode, "inherit");
let (sink, _rx) = EventSink::channel();
let second = exec
.run(
run_spec_with_messages(
"what nonce?",
vec![
message("user", "remember amber-572"),
message("assistant", "stored"),
message("user", "what nonce?"),
],
),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(second.status, TerminalStatus::Completed);
let argv = std::fs::read_to_string(bin_dir.path().join("argv-2.txt")).unwrap();
let arguments = argv.lines().collect::<Vec<_>>();
let resume_index = arguments
.iter()
.position(|argument| *argument == "resume")
.expect("resume subcommand");
assert_eq!(arguments.get(resume_index + 1), Some(&"thread-one"));
assert_eq!(arguments.last(), Some(&"-"));
assert_eq!(
std::fs::read_to_string(bin_dir.path().join("stdin-2.txt"))
.unwrap()
.trim(),
"what nonce?"
);
let state: CodexSessionState =
serde_json::from_slice(&std::fs::read(&state_path).unwrap()).unwrap();
assert_eq!(state.thread_id, "thread-two");
let (sink, _rx) = EventSink::channel();
let third = exec
.run(
run_spec("fresh rerun"),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(third.status, TerminalStatus::Completed);
assert!(!state_path.exists());
assert!(!std::fs::read_to_string(bin_dir.path().join("argv-3.txt"))
.unwrap()
.lines()
.any(|argument| argument == "resume"));
}
#[tokio::test]
async fn missing_or_mismatched_state_uses_bounded_history_rehydration() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let state_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
printf '%s\n' "$@" > "$DIR/argv.txt"
cat > "$DIR/stdin.txt"
echo '{"type":"thread.started","thread_id":"fallback-thread"}'
echo '{"type":"turn.started"}'
echo '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"ok"}}'
echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
"#,
);
let mut exec = executor(bin, workspace.path());
exec.state_dir = Some(state_dir.path().to_path_buf());
write_json_atomic(
&state_dir.path().join(CODEX_SESSION_STATE_FILE),
&CodexSessionState {
thread_id: "unusable-thread".to_string(),
workspace: exec.workspace.clone(),
codex_home_mode: "isolated".to_string(),
updated_at: Utc::now(),
},
)
.await
.unwrap();
let (sink, _rx) = EventSink::channel();
let outcome = exec
.run(
run_spec_with_messages(
"continue",
vec![
message("user", "old question"),
message("assistant", "old answer"),
message("user", "continue"),
],
),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Completed);
let argv = std::fs::read_to_string(bin_dir.path().join("argv.txt")).unwrap();
assert!(!argv.lines().any(|argument| argument == "resume"));
assert!(!argv.contains("unusable-thread"));
let body = std::fs::read_to_string(bin_dir.path().join("stdin.txt")).unwrap();
assert!(body.contains("## Prior conversation (rehydrated)"));
assert!(body.contains("old question"));
assert!(body.contains("old answer"));
assert!(body.contains("## Current task"));
assert_eq!(body.matches("continue").count(), 1);
let state: CodexSessionState = serde_json::from_slice(
&std::fs::read(state_dir.path().join(CODEX_SESSION_STATE_FILE)).unwrap(),
)
.unwrap();
assert_eq!(state.thread_id, "fallback-thread");
let workspace_mismatch = CodexSessionState {
thread_id: "wrong-workspace".to_string(),
workspace: Some("/different/workspace".to_string()),
codex_home_mode: "inherit".to_string(),
updated_at: Utc::now(),
};
write_json_atomic(
&state_dir.path().join(CODEX_SESSION_STATE_FILE),
&workspace_mismatch,
)
.await
.unwrap();
assert_eq!(exec.resolve_resume_id().await, None);
}
#[tokio::test]
async fn invalid_resume_retries_once_fresh_with_rehydrated_history() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let state_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
N=$(cat "$DIR/count" 2>/dev/null || echo 0)
N=$((N+1))
echo "$N" > "$DIR/count"
printf '%s\n' "$@" > "$DIR/argv-$N.txt"
cat > "$DIR/stdin-$N.txt"
case " $* " in
*' resume dead-thread '*)
echo 'thread not found' >&2
exit 1
;;
*)
echo '{"type":"thread.started","thread_id":"recovered-thread"}'
echo '{"type":"turn.started"}'
echo '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"recovered"}}'
echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
;;
esac
"#,
);
let mut exec = executor(bin, workspace.path());
exec.state_dir = Some(state_dir.path().to_path_buf());
write_json_atomic(
&state_dir.path().join(CODEX_SESSION_STATE_FILE),
&CodexSessionState {
thread_id: "dead-thread".to_string(),
workspace: exec.workspace.clone(),
codex_home_mode: "inherit".to_string(),
updated_at: Utc::now(),
},
)
.await
.unwrap();
let (sink, mut events) = EventSink::channel();
let outcome = exec
.run(
run_spec_with_messages(
"please continue",
vec![
message("user", "prior nonce amber-572"),
message("assistant", "stored"),
message("user", "please continue"),
],
),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Completed);
assert_eq!(outcome.result.as_deref(), Some("recovered"));
assert_eq!(
std::fs::read_to_string(bin_dir.path().join("count"))
.unwrap()
.trim(),
"2"
);
let first_argv = std::fs::read_to_string(bin_dir.path().join("argv-1.txt")).unwrap();
assert!(first_argv.lines().any(|argument| argument == "resume"));
assert!(first_argv.lines().any(|argument| argument == "dead-thread"));
let second_argv = std::fs::read_to_string(bin_dir.path().join("argv-2.txt")).unwrap();
assert!(!second_argv.lines().any(|argument| argument == "resume"));
let fallback = std::fs::read_to_string(bin_dir.path().join("stdin-2.txt")).unwrap();
assert!(fallback.contains("## Prior conversation (rehydrated)"));
assert!(fallback.contains("prior nonce amber-572"));
assert!(std::iter::from_fn(|| events.try_recv().ok()).any(|event| {
event["type"] == "runner_progress" && event["phase"] == "resume_fallback"
}));
let state: CodexSessionState = serde_json::from_slice(
&std::fs::read(state_dir.path().join(CODEX_SESSION_STATE_FILE)).unwrap(),
)
.unwrap();
assert_eq!(state.thread_id, "recovered-thread");
assert!(!std::fs::read_dir(state_dir.path())
.unwrap()
.filter_map(Result::ok)
.any(|entry| entry.file_name().to_string_lossy().contains(".tmp.")));
}
#[tokio::test]
async fn resume_failure_after_turn_progress_is_not_retried() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let state_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
N=$(cat "$DIR/count" 2>/dev/null || echo 0)
N=$((N+1))
echo "$N" > "$DIR/count"
cat >/dev/null
echo '{"type":"thread.started","thread_id":"progressed-thread"}'
echo '{"type":"turn.started"}'
echo '{"type":"turn.failed","error":{"message":"model failed"}}'
exit 1
"#,
);
let mut exec = executor(bin, workspace.path());
exec.state_dir = Some(state_dir.path().to_path_buf());
write_json_atomic(
&state_dir.path().join(CODEX_SESSION_STATE_FILE),
&CodexSessionState {
thread_id: "resumable-thread".to_string(),
workspace: exec.workspace.clone(),
codex_home_mode: "inherit".to_string(),
updated_at: Utc::now(),
},
)
.await
.unwrap();
let (sink, mut events) = EventSink::channel();
let outcome = exec
.run(
run_spec_with_messages(
"continue",
vec![message("user", "earlier"), message("user", "continue")],
),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Error);
assert_eq!(
std::fs::read_to_string(bin_dir.path().join("count"))
.unwrap()
.trim(),
"1"
);
assert!(!std::iter::from_fn(|| events.try_recv().ok()).any(|event| {
event["type"] == "runner_progress" && event["phase"] == "resume_fallback"
}));
}
#[tokio::test]
async fn failed_fallback_is_not_retried_a_second_time() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let state_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
N=$(cat "$DIR/count" 2>/dev/null || echo 0)
N=$((N+1))
echo "$N" > "$DIR/count"
cat >/dev/null
echo "attempt $N failed" >&2
exit 1
"#,
);
let mut exec = executor(bin, workspace.path());
exec.state_dir = Some(state_dir.path().to_path_buf());
let state_path = state_dir.path().join(CODEX_SESSION_STATE_FILE);
write_json_atomic(
&state_path,
&CodexSessionState {
thread_id: "dead-thread".to_string(),
workspace: exec.workspace.clone(),
codex_home_mode: "inherit".to_string(),
updated_at: Utc::now(),
},
)
.await
.unwrap();
let (sink, mut events) = EventSink::channel();
let outcome = exec
.run(
run_spec_with_messages(
"continue",
vec![message("user", "earlier"), message("user", "continue")],
),
sink,
SteerInbox::disconnected(),
CancellationToken::new(),
)
.await;
assert_eq!(outcome.status, TerminalStatus::Error);
assert_eq!(
std::fs::read_to_string(bin_dir.path().join("count"))
.unwrap()
.trim(),
"2"
);
assert_eq!(
std::iter::from_fn(|| events.try_recv().ok())
.filter(|event| {
event["type"] == "runner_progress" && event["phase"] == "resume_fallback"
})
.count(),
1
);
assert!(!state_path.exists());
}
#[tokio::test]
async fn cancellation_kills_the_entire_process_group() {
let workspace = tempfile::tempdir().unwrap();
let bin_dir = tempfile::tempdir().unwrap();
let bin = write_stub(
bin_dir.path(),
r#"
DIR=$(cd "$(dirname "$0")" && pwd)
read -r prompt
echo '{"type":"thread.started","thread_id":"cancel-thread"}'
(trap '' TERM; sleep 30) &
echo $! > "$DIR/grandchild.pid"
wait
"#,
);
let pid_path = bin_dir.path().join("grandchild.pid");
let (sink, _rx) = EventSink::channel();
let cancel = CancellationToken::new();
let cancel_for_run = cancel.clone();
let run = tokio::spawn(async move {
executor(bin, workspace.path())
.run(
run_spec("wait"),
sink,
SteerInbox::disconnected(),
cancel_for_run,
)
.await
});
for _ in 0..500 {
if pid_path.exists() {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(pid_path.exists(), "grandchild pid was recorded");
let grandchild_pid: libc::pid_t = std::fs::read_to_string(&pid_path)
.unwrap()
.trim()
.parse()
.unwrap();
cancel.cancel();
let outcome = tokio::time::timeout(Duration::from_secs(10), run)
.await
.expect("cancel completed within TERM/KILL bound")
.unwrap();
assert_eq!(outcome.status, TerminalStatus::Cancelled);
for _ in 0..100 {
if unsafe { libc::kill(grandchild_pid, 0) } == -1 {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(unsafe { libc::kill(grandchild_pid, 0) }, -1);
}
}
}