Skip to main content

bamboo_agent/
codex_cli_executor.rs

1//! `CodexExecutor`: a [`ChildExecutor`] that drives the official OpenAI
2//! Codex CLI through `codex exec --json`.
3//!
4//! The CLI is one process per activation. Prompts are written on stdin (never
5//! argv), stdout is consumed as bounded JSONL, and the process owns a process
6//! group so cancellation tears down any descendants as well as the leader.
7//! Provider/auth selection and Bamboo permission-profile mapping are resolved
8//! before every spawn. Thread ids are persisted per child so later activations
9//! can use `codex exec resume`, with bounded history rehydration when the local
10//! Codex transcript is unavailable.
11
12use std::collections::{HashMap, HashSet};
13use std::path::{Path, PathBuf};
14use std::process::Stdio;
15use std::sync::Arc;
16use std::time::Duration;
17
18use async_trait::async_trait;
19use chrono::{DateTime, Utc};
20use serde::{Deserialize, Serialize};
21use serde_json::{json, Value};
22use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
23use tokio::process::{Child, Command};
24use tokio::sync::Mutex;
25use tokio_util::sync::CancellationToken;
26
27use bamboo_agent_core::{AgentEvent, TokenUsage, ToolResult};
28use bamboo_subagent::codex_discovery::discover_codex_cli;
29use bamboo_subagent::executor::{ChildExecutor, ChildOutcome, EventSink, SteerInbox};
30use bamboo_subagent::executor_util::{build_rehydrated_turn, write_json_atomic};
31use bamboo_subagent::proto::RunSpec;
32
33/// The oldest Codex CLI schema this executor intentionally supports. The
34/// executor additionally capability-checks `exec --help` and `exec resume
35/// --help`, so a backported or vendor build must still expose the required
36/// flags. Version 0.144 is the schema verified by issue #569.
37pub use bamboo_subagent::codex_discovery::MIN_CODEX_VERSION;
38
39const MAX_STDOUT_LINE_BYTES: usize = 10 * 1024 * 1024;
40const STDERR_TAIL_BYTES: usize = 16 * 1024;
41const TOOL_RESULT_TRUNCATE_CHARS: usize = 20_000;
42const SIGTERM_WAIT: Duration = Duration::from_secs(5);
43const PROCESS_EXIT_WAIT: Duration = Duration::from_secs(5);
44
45const ENV_ALLOWLIST: &[&str] = &[
46    "HOME", "PATH", "SHELL", "TERM", "LANG", "TMPDIR", "USER", "LOGNAME",
47];
48
49const CODEX_PROVIDER_ENV: &str = "BAMBOO_CODEX_PROVIDER_KEY";
50const CODEX_SESSION_STATE_FILE: &str = "codex-session.json";
51
52#[derive(Debug, Clone, Serialize, Deserialize)]
53struct CodexSessionState {
54    thread_id: String,
55    workspace: Option<String>,
56    codex_home_mode: String,
57    updated_at: DateTime<Utc>,
58}
59
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum CodexAuthMode {
62    Inherit,
63    ApiKey,
64    Custom,
65    Bamboo,
66}
67
68impl CodexAuthMode {
69    pub(crate) fn as_str(self) -> &'static str {
70        match self {
71            Self::Inherit => "inherit",
72            Self::ApiKey => "api_key",
73            Self::Custom => "custom",
74            Self::Bamboo => "bamboo",
75        }
76    }
77}
78
79/// Fully resolved auth posture. The custom-provider key is deliberately
80/// private and this type has no `Debug` implementation, preventing accidental
81/// secret logging from executor/spec diagnostics.
82#[derive(Clone)]
83pub struct CodexAuthConfig {
84    mode: CodexAuthMode,
85    base_url: Option<String>,
86    wire_api: String,
87    provider_key: Option<String>,
88}
89
90impl CodexAuthConfig {
91    pub fn inherit() -> Self {
92        Self {
93            mode: CodexAuthMode::Inherit,
94            base_url: None,
95            wire_api: "responses".to_string(),
96            provider_key: None,
97        }
98    }
99
100    pub fn mode(&self) -> CodexAuthMode {
101        self.mode
102    }
103
104    pub(crate) fn isolated(&self) -> bool {
105        self.mode != CodexAuthMode::Inherit
106    }
107
108    fn generated_config_toml(&self) -> Result<String, String> {
109        if !self.isolated() || self.mode == CodexAuthMode::ApiKey {
110            return Ok("# Generated by Bamboo; authentication is environment-only.\n".to_string());
111        }
112        let provider_id = match self.mode {
113            CodexAuthMode::Custom => "custom",
114            CodexAuthMode::Bamboo => "bamboo",
115            CodexAuthMode::Inherit | CodexAuthMode::ApiKey => unreachable!(),
116        };
117        let base_url = self.base_url.as_ref().ok_or_else(|| {
118            format!("Codex auth mode '{provider_id}' requires a provider base URL")
119        })?;
120        let mut provider = toml::Table::new();
121        provider.insert(
122            "name".to_string(),
123            toml::Value::String(format!("Bamboo {provider_id}")),
124        );
125        provider.insert(
126            "base_url".to_string(),
127            toml::Value::String(base_url.clone()),
128        );
129        provider.insert(
130            "env_key".to_string(),
131            toml::Value::String(CODEX_PROVIDER_ENV.to_string()),
132        );
133        provider.insert(
134            "wire_api".to_string(),
135            toml::Value::String(self.wire_api.clone()),
136        );
137        let mut providers = toml::Table::new();
138        providers.insert(provider_id.to_string(), toml::Value::Table(provider));
139        let mut root = toml::Table::new();
140        root.insert(
141            "model_provider".to_string(),
142            toml::Value::String(provider_id.to_string()),
143        );
144        root.insert("model_providers".to_string(), toml::Value::Table(providers));
145        toml::to_string(&toml::Value::Table(root))
146            .map_err(|error| format!("serialize isolated Codex config.toml: {error}"))
147    }
148
149    /// App-server stays alive across activations, while Bamboo provider tokens
150    /// are intentionally minted and revoked per activation. Use Codex's
151    /// command-backed auth hook to read the current 0600 token file instead of
152    /// freezing the first token in the process environment.
153    pub(crate) fn generated_app_server_config_toml(
154        &self,
155        token_helper: &Path,
156        token_path: &Path,
157    ) -> Result<String, String> {
158        if self.mode != CodexAuthMode::Bamboo {
159            return self.generated_config_toml();
160        }
161        let base_url = self
162            .base_url
163            .as_ref()
164            .ok_or_else(|| "Codex auth mode 'bamboo' requires a provider base URL".to_string())?;
165        let mut auth = toml::Table::new();
166        auth.insert(
167            "command".to_string(),
168            toml::Value::String(token_helper.to_string_lossy().into_owned()),
169        );
170        auth.insert(
171            "args".to_string(),
172            toml::Value::Array(vec![
173                toml::Value::String("codex-provider-token".to_string()),
174                toml::Value::String(token_path.to_string_lossy().into_owned()),
175            ]),
176        );
177        auth.insert("timeout_ms".to_string(), toml::Value::Integer(5_000));
178        auth.insert("refresh_interval_ms".to_string(), toml::Value::Integer(1));
179        let mut provider = toml::Table::new();
180        provider.insert(
181            "name".to_string(),
182            toml::Value::String("Bamboo bamboo".to_string()),
183        );
184        provider.insert(
185            "base_url".to_string(),
186            toml::Value::String(base_url.clone()),
187        );
188        provider.insert(
189            "wire_api".to_string(),
190            toml::Value::String(self.wire_api.clone()),
191        );
192        provider.insert("auth".to_string(), toml::Value::Table(auth));
193        let mut providers = toml::Table::new();
194        providers.insert("bamboo".to_string(), toml::Value::Table(provider));
195        let mut root = toml::Table::new();
196        root.insert(
197            "model_provider".to_string(),
198            toml::Value::String("bamboo".to_string()),
199        );
200        root.insert("model_providers".to_string(), toml::Value::Table(providers));
201        toml::to_string(&toml::Value::Table(root))
202            .map_err(|error| format!("serialize isolated Codex app-server config.toml: {error}"))
203    }
204
205    pub(crate) fn provider_key(&self) -> Option<&str> {
206        self.provider_key.as_deref()
207    }
208}
209
210/// Read the short-lived Bamboo provider token for Codex command-backed auth.
211/// The final path component is opened without following symlinks on Unix, and
212/// permissions are verified on the opened descriptor to avoid check/open races.
213pub fn read_codex_provider_token(path: &Path) -> Result<String, String> {
214    #[cfg(not(unix))]
215    {
216        let metadata = std::fs::symlink_metadata(path)
217            .map_err(|error| format!("inspect Codex provider token: {error}"))?;
218        if metadata.file_type().is_symlink() {
219            return Err("Codex provider token path must not be a symlink".to_string());
220        }
221    }
222
223    let mut options = std::fs::OpenOptions::new();
224    options.read(true);
225    #[cfg(unix)]
226    {
227        use std::os::unix::fs::OpenOptionsExt as _;
228        options.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC);
229    }
230    let mut file = options
231        .open(path)
232        .map_err(|error| format!("open Codex provider token: {error}"))?;
233    let metadata = file
234        .metadata()
235        .map_err(|error| format!("inspect Codex provider token: {error}"))?;
236    if !metadata.is_file() {
237        return Err("Codex provider token path must be a regular file".to_string());
238    }
239    #[cfg(unix)]
240    {
241        use std::os::unix::fs::PermissionsExt as _;
242        if metadata.permissions().mode() & 0o077 != 0 {
243            return Err(
244                "Codex provider token file must not be accessible by group/other".to_string(),
245            );
246        }
247    }
248    let mut token = String::new();
249    std::io::Read::read_to_string(&mut file, &mut token)
250        .map_err(|error| format!("read Codex provider token: {error}"))?;
251    let token = token.trim();
252    if token.is_empty() {
253        return Err("Codex provider token file is empty".to_string());
254    }
255    Ok(token.to_string())
256}
257
258/// Resolve and validate the public provision fields plus the one referenced
259/// provider credential carried in the in-memory secrets envelope.
260#[allow(clippy::too_many_arguments)]
261pub fn resolve_codex_auth_config(
262    auth_mode: Option<&str>,
263    legacy_inherit_user_config: bool,
264    base_url: Option<String>,
265    wire_api: Option<String>,
266    provider_key_ref: Option<&str>,
267    credentials: &[bamboo_subagent::provision::ScopedCredential],
268    forward_env: &[String],
269) -> Result<CodexAuthConfig, String> {
270    let mode = match auth_mode {
271        Some("inherit") => CodexAuthMode::Inherit,
272        Some("api_key") => CodexAuthMode::ApiKey,
273        Some("custom") => CodexAuthMode::Custom,
274        Some("bamboo") => CodexAuthMode::Bamboo,
275        Some(other) => {
276            return Err(format!(
277                "unknown Codex auth mode '{other}'; expected inherit, api_key, custom, or bamboo"
278            ))
279        }
280        None if legacy_inherit_user_config => CodexAuthMode::Inherit,
281        None => CodexAuthMode::Bamboo,
282    };
283    let wire_api = wire_api.unwrap_or_else(|| "responses".to_string());
284    if wire_api != "responses" {
285        return Err(format!(
286            "unsupported Codex wire_api '{wire_api}'; Codex CLI >= 0.144 requires responses"
287        ));
288    }
289    validate_codex_base_url(mode, base_url.as_deref())?;
290    validate_codex_forward_env(mode, forward_env)?;
291
292    let provider_key = if mode == CodexAuthMode::Custom {
293        let reference = provider_key_ref.ok_or_else(|| {
294            "Codex auth mode 'custom' requires codex_provider_key_ref".to_string()
295        })?;
296        Some(
297            credentials
298                .iter()
299                .find(|credential| credential.credential_ref.as_deref() == Some(reference))
300                .map(|credential| credential.api_key.clone())
301                .ok_or_else(|| {
302                    format!(
303                        "Codex custom provider credential reference '{reference}' did not resolve"
304                    )
305                })?,
306        )
307    } else {
308        None
309    };
310
311    Ok(CodexAuthConfig {
312        mode,
313        base_url,
314        wire_api,
315        provider_key,
316    })
317}
318
319/// Resolve the per-child directory used for `--output-last-message`.
320pub fn resolve_codex_state_dir(storage_dir: &Option<String>, child_id: &str) -> PathBuf {
321    storage_dir
322        .clone()
323        .map(PathBuf::from)
324        .unwrap_or_else(|| bamboo_config::paths::subagents_dir().join(child_id))
325}
326
327#[derive(Debug, Clone, Copy, PartialEq, Eq)]
328enum CodexSandbox {
329    ReadOnly,
330    WorkspaceWrite,
331    DangerFullAccess,
332}
333
334impl CodexSandbox {
335    fn parse(raw: &str) -> Result<Self, String> {
336        match raw {
337            "read-only" => Ok(Self::ReadOnly),
338            "workspace-write" => Ok(Self::WorkspaceWrite),
339            "danger-full-access" => Ok(Self::DangerFullAccess),
340            other => Err(format!(
341                "unknown Codex sandbox '{other}'; expected read-only, workspace-write, or danger-full-access"
342            )),
343        }
344    }
345
346    fn as_str(self) -> &'static str {
347        match self {
348            Self::ReadOnly => "read-only",
349            Self::WorkspaceWrite => "workspace-write",
350            Self::DangerFullAccess => "danger-full-access",
351        }
352    }
353}
354
355#[derive(Debug, Clone, Copy, PartialEq, Eq)]
356enum CodexApprovalPolicy {
357    Never,
358    OnFailure,
359    OnRequest,
360}
361
362impl CodexApprovalPolicy {
363    fn parse(raw: &str) -> Result<Self, String> {
364        match raw {
365            "never" => Ok(Self::Never),
366            "on-failure" => Ok(Self::OnFailure),
367            "on-request" => Ok(Self::OnRequest),
368            "untrusted" => Err("Codex approval policy 'untrusted' is unsupported; use on-request in app_server mode".to_string()),
369            other => Err(format!(
370                "unknown Codex approval policy '{other}'; expected never, on-failure, or on-request"
371            )),
372        }
373    }
374
375    fn as_str(self) -> &'static str {
376        match self {
377            Self::Never => "never",
378            Self::OnFailure => "on-failure",
379            Self::OnRequest => "on-request",
380        }
381    }
382}
383
384#[derive(Debug, Clone, Copy, PartialEq, Eq)]
385enum CodexPolicyInvocation {
386    Explicit,
387    FullAuto,
388    DangerBypass,
389}
390
391impl CodexPolicyInvocation {
392    fn as_str(self) -> &'static str {
393        match self {
394            Self::Explicit => "explicit",
395            Self::FullAuto => "full-auto",
396            Self::DangerBypass => "danger-bypass",
397        }
398    }
399}
400
401#[derive(Debug, Clone)]
402pub struct CodexPermissionConfig {
403    sandbox: Option<CodexSandbox>,
404    approval_policy: Option<CodexApprovalPolicy>,
405    network_access: bool,
406    allow_danger_bypass: bool,
407    permission_profile: Option<String>,
408    provisioned_permission: bamboo_domain::PermissionModeResolution,
409    provisioned_session_id: Option<String>,
410    workspace_owned: bool,
411}
412
413#[derive(Debug, Clone, PartialEq, Eq)]
414struct EffectiveCodexPolicy {
415    sandbox: CodexSandbox,
416    approval_policy: CodexApprovalPolicy,
417    invocation: CodexPolicyInvocation,
418    network_access: bool,
419    warnings: Vec<String>,
420}
421
422impl CodexPermissionConfig {
423    fn effective(&self, bypass: bool, is_root: bool) -> EffectiveCodexPolicy {
424        let read_only_profile = self
425            .permission_profile
426            .as_deref()
427            .is_some_and(profile_is_read_only);
428        let mut warnings = Vec::new();
429
430        let (sandbox, approval_policy, invocation) = match self.sandbox {
431            Some(CodexSandbox::ReadOnly) => (
432                CodexSandbox::ReadOnly,
433                self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
434                CodexPolicyInvocation::Explicit,
435            ),
436            Some(CodexSandbox::WorkspaceWrite) => (
437                CodexSandbox::WorkspaceWrite,
438                self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
439                CodexPolicyInvocation::Explicit,
440            ),
441            Some(CodexSandbox::DangerFullAccess) => {
442                resolve_danger_request(bypass, self.allow_danger_bypass, is_root, &mut warnings)
443            }
444            None if self.allow_danger_bypass => {
445                resolve_danger_request(bypass, true, is_root, &mut warnings)
446            }
447            None if bypass => (
448                CodexSandbox::WorkspaceWrite,
449                CodexApprovalPolicy::Never,
450                CodexPolicyInvocation::FullAuto,
451            ),
452            None if read_only_profile => (
453                CodexSandbox::ReadOnly,
454                self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
455                CodexPolicyInvocation::Explicit,
456            ),
457            None => (
458                CodexSandbox::WorkspaceWrite,
459                self.approval_policy.unwrap_or(CodexApprovalPolicy::Never),
460                CodexPolicyInvocation::Explicit,
461            ),
462        };
463
464        let network_access = match sandbox {
465            CodexSandbox::ReadOnly => false,
466            CodexSandbox::WorkspaceWrite => self.network_access,
467            CodexSandbox::DangerFullAccess => true,
468        };
469        if invocation == CodexPolicyInvocation::DangerBypass {
470            warnings.push(
471                "DANGER: Codex is running with approvals and the OS sandbox disabled for this parent-bypass run"
472                    .to_string(),
473            );
474        }
475
476        EffectiveCodexPolicy {
477            sandbox,
478            approval_policy,
479            invocation,
480            network_access,
481            warnings,
482        }
483    }
484
485    /// Resolve the executor posture for one session activation.
486    ///
487    /// Auto removes only the approval gate: it forces Codex's supported
488    /// `never` policy without widening the configured sandbox or borrowing the
489    /// legacy bypass path that can opt into danger-full-access.
490    fn effective_for_session(
491        &self,
492        permission: bamboo_domain::PermissionModeResolution,
493        is_root: bool,
494    ) -> EffectiveCodexPolicy {
495        let mut policy = self.effective(permission.bypass_permissions(), is_root);
496        if permission.effective == bamboo_domain::PermissionMode::Plan {
497            policy.sandbox = CodexSandbox::ReadOnly;
498            policy.approval_policy = CodexApprovalPolicy::Never;
499            policy.invocation = CodexPolicyInvocation::Explicit;
500            policy.network_access = false;
501        } else if permission.suppress_approval_prompts() {
502            policy.approval_policy = CodexApprovalPolicy::Never;
503        }
504        policy
505    }
506}
507
508fn resolve_danger_request(
509    bypass: bool,
510    allow_danger_bypass: bool,
511    is_root: bool,
512    warnings: &mut Vec<String>,
513) -> (CodexSandbox, CodexApprovalPolicy, CodexPolicyInvocation) {
514    if bypass && allow_danger_bypass && !is_root {
515        return (
516            CodexSandbox::DangerFullAccess,
517            CodexApprovalPolicy::Never,
518            CodexPolicyInvocation::DangerBypass,
519        );
520    }
521
522    let reason = if !bypass {
523        "the parent run is not in bypass mode"
524    } else if !allow_danger_bypass {
525        "codex_allow_danger_bypass is false"
526    } else {
527        "the worker is running as root"
528    };
529    warnings.push(format!(
530        "Codex danger-full-access request was downgraded to a sandboxed policy because {reason}"
531    ));
532    // A requested danger posture first falls back to Codex's sandboxed
533    // convenience mode. This keeps the two gates independent: neither a
534    // config-only request nor a root process can silently become unsandboxed.
535    (
536        CodexSandbox::WorkspaceWrite,
537        CodexApprovalPolicy::Never,
538        CodexPolicyInvocation::FullAuto,
539    )
540}
541
542fn profile_is_read_only(raw: &str) -> bool {
543    let normalized = raw.trim().to_ascii_lowercase().replace(['_', ' '], "-");
544    matches!(
545        normalized.as_str(),
546        "read-only" | "readonly" | "research" | "researcher" | "guardian" | "plan"
547    )
548}
549
550#[allow(clippy::too_many_arguments)]
551pub fn resolve_codex_permission_config(
552    sandbox: Option<&str>,
553    approval_policy: Option<&str>,
554    network_access: bool,
555    allow_danger_bypass: bool,
556    permission_profile: Option<String>,
557    provisioned_bypass: bool,
558    workspace_owned: bool,
559) -> Result<CodexPermissionConfig, String> {
560    let sandbox = sandbox.map(CodexSandbox::parse).transpose()?;
561    let approval_policy = approval_policy
562        .map(CodexApprovalPolicy::parse)
563        .transpose()?;
564    if approval_policy == Some(CodexApprovalPolicy::OnRequest) {
565        return Err(
566            "Codex approval policy 'on-request' requires codex_mode = \"app_server\"; non-interactive exec mode has no approval relay"
567                .to_string(),
568        );
569    }
570    let profile_read_only = permission_profile
571        .as_deref()
572        .is_some_and(profile_is_read_only);
573    if network_access
574        && (sandbox == Some(CodexSandbox::ReadOnly) || (sandbox.is_none() && profile_read_only))
575    {
576        return Err(
577            "Codex network access requires an effective workspace-write sandbox".to_string(),
578        );
579    }
580    Ok(CodexPermissionConfig {
581        sandbox,
582        approval_policy,
583        network_access,
584        allow_danger_bypass,
585        permission_profile,
586        provisioned_permission: bamboo_domain::resolve_permission_mode(
587            if provisioned_bypass {
588                bamboo_domain::SessionPermissionMode::Bypass
589            } else {
590                bamboo_domain::SessionPermissionMode::Default
591            },
592            bamboo_domain::PermissionMode::Default,
593        ),
594        provisioned_session_id: None,
595        workspace_owned,
596    })
597}
598
599/// App-server mode always routes approvals to Bamboo. Accepting `never` or
600/// `on-failure` here would make the selected transport's safety contract lie.
601#[allow(clippy::too_many_arguments)]
602pub fn resolve_codex_app_server_permission_config(
603    sandbox: Option<&str>,
604    approval_policy: Option<&str>,
605    network_access: bool,
606    allow_danger_bypass: bool,
607    permission_profile: Option<String>,
608    provisioned_bypass: bool,
609    workspace_owned: bool,
610) -> Result<CodexPermissionConfig, String> {
611    match approval_policy {
612        None | Some("on-request") => {}
613        Some(other) => {
614            return Err(format!(
615                "Codex approval policy '{other}' is incompatible with codex_mode = \"app_server\"; use on-request"
616            ))
617        }
618    }
619    let sandbox = sandbox.map(CodexSandbox::parse).transpose()?;
620    let profile_read_only = permission_profile
621        .as_deref()
622        .is_some_and(profile_is_read_only);
623    if network_access
624        && (sandbox == Some(CodexSandbox::ReadOnly) || (sandbox.is_none() && profile_read_only))
625    {
626        return Err(
627            "Codex network access requires an effective workspace-write sandbox".to_string(),
628        );
629    }
630    Ok(CodexPermissionConfig {
631        sandbox,
632        approval_policy: Some(CodexApprovalPolicy::OnRequest),
633        network_access,
634        allow_danger_bypass,
635        permission_profile,
636        provisioned_permission: bamboo_domain::resolve_permission_mode(
637            if provisioned_bypass {
638                bamboo_domain::SessionPermissionMode::Bypass
639            } else {
640                bamboo_domain::SessionPermissionMode::Default
641            },
642            bamboo_domain::PermissionMode::Default,
643        ),
644        provisioned_session_id: None,
645        workspace_owned,
646    })
647}
648
649impl CodexPermissionConfig {
650    pub(crate) fn with_provisioned_permission_resolution(
651        mut self,
652        resolution: bamboo_domain::PermissionModeResolution,
653        session_id: String,
654    ) -> Self {
655        self.provisioned_permission = resolution;
656        self.provisioned_session_id = Some(session_id);
657        self
658    }
659
660    pub(crate) fn app_server_posture(
661        &self,
662        permission: bamboo_domain::PermissionModeResolution,
663    ) -> (String, bool, Vec<String>) {
664        let mut policy = self.effective(permission.bypass_permissions(), running_as_root());
665        if permission.effective == bamboo_domain::PermissionMode::Plan {
666            policy.sandbox = CodexSandbox::ReadOnly;
667            policy.network_access = false;
668            policy.invocation = CodexPolicyInvocation::Explicit;
669        }
670        policy.approval_policy = CodexApprovalPolicy::OnRequest;
671        (
672            policy.sandbox.as_str().to_string(),
673            policy.network_access,
674            policy.warnings,
675        )
676    }
677
678    pub(crate) fn permission_profile(&self) -> Option<&str> {
679        self.permission_profile.as_deref()
680    }
681
682    pub(crate) fn provisioned_session_id(&self) -> Option<&str> {
683        self.provisioned_session_id.as_deref()
684    }
685
686    pub(crate) fn activation_permission_resolution(
687        &self,
688        context: Option<&bamboo_subagent::proto::PermissionPolicyContext>,
689    ) -> Result<bamboo_domain::PermissionModeResolution, String> {
690        let Some(context) = context else {
691            return Ok(self.provisioned_permission);
692        };
693        let (requested, effective) = context.resolved_modes()?;
694        Ok(bamboo_domain::PermissionModeResolution {
695            requested,
696            effective,
697        })
698    }
699
700    pub(crate) fn activation_permission(
701        &self,
702        context: Option<&bamboo_subagent::proto::PermissionPolicyContext>,
703    ) -> Result<ResolvedCodexPermission, String> {
704        let resolution = self.activation_permission_resolution(context)?;
705        let Some(context) = context else {
706            return Ok(ResolvedCodexPermission {
707                resolution,
708                policy_revision: 0,
709                explicit_deny_reason: None,
710            });
711        };
712        let policy =
713            serde_json::from_value::<bamboo_tools::permission::SerializablePermissionConfig>(
714                context.policy.clone(),
715            )
716            .map_err(|error| format!("decode Codex permission policy: {error}"))?;
717        Ok(ResolvedCodexPermission {
718            resolution,
719            policy_revision: context.revision,
720            explicit_deny_reason: bamboo_tools::permission::explicit_deny_policy_reason(&policy),
721        })
722    }
723}
724
725#[derive(Debug, Clone)]
726pub(crate) struct ResolvedCodexPermission {
727    pub resolution: bamboo_domain::PermissionModeResolution,
728    pub policy_revision: u64,
729    pub explicit_deny_reason: Option<String>,
730}
731
732#[cfg(unix)]
733fn running_as_root() -> bool {
734    // SAFETY: geteuid has no preconditions and does not dereference memory.
735    unsafe { libc::geteuid() == 0 }
736}
737
738#[cfg(not(unix))]
739fn running_as_root() -> bool {
740    false
741}
742
743/// One-process-per-activation Codex CLI executor.
744pub struct CodexExecutor {
745    binary: PathBuf,
746    version: String,
747    model: Option<String>,
748    permissions: CodexPermissionConfig,
749    workspace: Option<String>,
750    state_dir: Option<PathBuf>,
751    forward_env: Vec<String>,
752    auth: CodexAuthConfig,
753}
754
755impl CodexExecutor {
756    /// Resolve and capability-check the configured Codex binary before the
757    /// worker begins serving runs. This keeps install/version failures at
758    /// provisioning time instead of surfacing halfway through a child turn.
759    #[allow(clippy::too_many_arguments)]
760    pub async fn new(
761        binary: Option<String>,
762        model: Option<String>,
763        workspace: Option<String>,
764        state_dir: Option<PathBuf>,
765        forward_env: Vec<String>,
766        auth: CodexAuthConfig,
767        permissions: CodexPermissionConfig,
768    ) -> Result<Self, String> {
769        let discovery = discover_codex_cli(binary.as_deref()).await?;
770
771        let executor = Self {
772            binary: PathBuf::from(discovery.path),
773            version: discovery.version,
774            model,
775            permissions,
776            workspace,
777            state_dir,
778            forward_env,
779            auth,
780        };
781        executor.prepare_auth_home().await?;
782        Ok(executor)
783    }
784
785    fn last_message_path(&self) -> Option<PathBuf> {
786        self.state_dir
787            .as_ref()
788            .map(|directory| directory.join("codex-last-message.txt"))
789    }
790
791    fn session_state_path(&self) -> Option<PathBuf> {
792        self.state_dir
793            .as_ref()
794            .map(|directory| directory.join(CODEX_SESSION_STATE_FILE))
795    }
796
797    fn codex_home_mode(&self) -> &'static str {
798        if self.auth.isolated() {
799            "isolated"
800        } else {
801            "inherit"
802        }
803    }
804
805    async fn read_session_state(&self) -> Option<CodexSessionState> {
806        let path = self.session_state_path()?;
807        let bytes = tokio::fs::read(path).await.ok()?;
808        serde_json::from_slice(&bytes).ok()
809    }
810
811    async fn delete_session_state(&self) {
812        let Some(path) = self.session_state_path() else {
813            return;
814        };
815        match tokio::fs::remove_file(&path).await {
816            Ok(()) => {}
817            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
818            Err(error) => tracing::warn!(
819                path = %path.display(),
820                %error,
821                "codex: remove stale session state"
822            ),
823        }
824    }
825
826    async fn write_session_state(&self, thread_id: &str) {
827        let Some(path) = self.session_state_path() else {
828            return;
829        };
830        let state = CodexSessionState {
831            thread_id: thread_id.to_string(),
832            workspace: self.workspace.clone(),
833            codex_home_mode: self.codex_home_mode().to_string(),
834            updated_at: Utc::now(),
835        };
836        if let Err(error) = write_json_atomic(&path, &state).await {
837            tracing::warn!(%error, "codex: persist thread state");
838        }
839    }
840
841    async fn resolve_resume_id(&self) -> Option<String> {
842        let state = self.read_session_state().await?;
843        if state.workspace != self.workspace {
844            tracing::warn!(
845                recorded = ?state.workspace,
846                current = ?self.workspace,
847                "codex: session state workspace mismatch; using history rehydration"
848            );
849            return None;
850        }
851        let current_mode = self.codex_home_mode();
852        if state.codex_home_mode != current_mode {
853            tracing::warn!(
854                recorded = %state.codex_home_mode,
855                current = current_mode,
856                "codex: session state CODEX_HOME mode mismatch; using history rehydration"
857            );
858            return None;
859        }
860        (!state.thread_id.trim().is_empty()).then_some(state.thread_id)
861    }
862
863    fn build_command(
864        &self,
865        run_provider_token: Option<&str>,
866        policy: &EffectiveCodexPolicy,
867        resume_id: Option<&str>,
868    ) -> Result<Command, String> {
869        let mut command = Command::new(&self.binary);
870        command
871            .arg("exec")
872            .arg("--json")
873            .arg("--color")
874            .arg("never");
875
876        let workspace_path = if let Some(workspace) = &self.workspace {
877            command.arg("--cd").arg(workspace);
878            Some(PathBuf::from(workspace))
879        } else {
880            std::env::current_dir().ok()
881        };
882        if self.permissions.workspace_owned
883            && workspace_path
884                .as_deref()
885                .is_some_and(|path| !has_git_metadata(path))
886        {
887            command.arg("--skip-git-repo-check");
888        }
889        if let Some(model) = &self.model {
890            command.arg("--model").arg(model);
891        }
892        match policy.invocation {
893            CodexPolicyInvocation::Explicit => {
894                command.arg("--sandbox").arg(policy.sandbox.as_str());
895                command.arg("--config").arg(format!(
896                    "approval_policy=\"{}\"",
897                    policy.approval_policy.as_str()
898                ));
899            }
900            CodexPolicyInvocation::FullAuto => {
901                command.arg("--full-auto");
902            }
903            CodexPolicyInvocation::DangerBypass => {
904                command.arg("--dangerously-bypass-approvals-and-sandbox");
905            }
906        }
907        if policy.sandbox == CodexSandbox::WorkspaceWrite && policy.network_access {
908            command
909                .arg("--config")
910                .arg("sandbox_workspace_write.network_access=true");
911        }
912        if self.auth.isolated() {
913            command.arg("--ignore-rules");
914        }
915        if let Some(path) = self.last_message_path() {
916            command.arg("--output-last-message").arg(path);
917        }
918
919        if let Some(thread_id) = resume_id {
920            command.arg("resume").arg(thread_id);
921        }
922
923        // `-` is the documented stdin prompt sentinel. It avoids both argv
924        // length limits and leaking the assignment through process listings.
925        command.arg("-");
926
927        command.env_clear();
928        for (key, value) in std::env::vars() {
929            if ENV_ALLOWLIST.contains(&key.as_str()) || key.starts_with("LC_") {
930                command.env(key, value);
931            }
932        }
933        if let Some(home) = self.codex_home() {
934            command.env("CODEX_HOME", home);
935        }
936        for name in &self.forward_env {
937            if let Ok(value) = std::env::var(name) {
938                command.env(name, value);
939            }
940        }
941        match self.auth.mode {
942            CodexAuthMode::Custom => {
943                let key = self.auth.provider_key.as_deref().ok_or_else(|| {
944                    "Codex custom provider key was not resolved at provisioning".to_string()
945                })?;
946                command.env(CODEX_PROVIDER_ENV, key);
947            }
948            CodexAuthMode::Bamboo => {
949                let token = run_provider_token.ok_or_else(|| {
950                    "Codex bamboo auth requires a per-run provider token".to_string()
951                })?;
952                command.env(CODEX_PROVIDER_ENV, token);
953            }
954            CodexAuthMode::Inherit | CodexAuthMode::ApiKey => {}
955        }
956        command.stdin(Stdio::piped());
957        command.stdout(Stdio::piped());
958        command.stderr(Stdio::piped());
959        command.kill_on_drop(true);
960        #[cfg(unix)]
961        command.process_group(0);
962        Ok(command)
963    }
964
965    fn codex_home(&self) -> Option<PathBuf> {
966        self.auth
967            .isolated()
968            .then(|| self.state_dir.as_ref().map(|path| path.join("codex-home")))
969            .flatten()
970    }
971
972    async fn prepare_auth_home(&self) -> Result<(), String> {
973        if !self.auth.isolated() {
974            return Ok(());
975        }
976        let home = self.codex_home().ok_or_else(|| {
977            "isolated Codex auth requires a Bamboo-managed state directory".to_string()
978        })?;
979        tokio::fs::create_dir_all(&home)
980            .await
981            .map_err(|error| format!("create isolated CODEX_HOME '{}': {error}", home.display()))?;
982        #[cfg(unix)]
983        tokio::fs::set_permissions(&home, std::os::unix::fs::PermissionsExt::from_mode(0o700))
984            .await
985            .map_err(|error| format!("secure isolated CODEX_HOME '{}': {error}", home.display()))?;
986
987        // A managed isolated home must never inherit a stale login artifact.
988        let auth_path = home.join("auth.json");
989        match tokio::fs::remove_file(&auth_path).await {
990            Ok(()) => {}
991            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
992            Err(error) => {
993                return Err(format!(
994                    "remove stale isolated Codex auth '{}': {error}",
995                    auth_path.display()
996                ))
997            }
998        }
999
1000        let config_path = home.join("config.toml");
1001        tokio::fs::write(&config_path, self.auth.generated_config_toml()?)
1002            .await
1003            .map_err(|error| {
1004                format!(
1005                    "write isolated Codex config '{}': {error}",
1006                    config_path.display()
1007                )
1008            })?;
1009        #[cfg(unix)]
1010        tokio::fs::set_permissions(
1011            &config_path,
1012            std::os::unix::fs::PermissionsExt::from_mode(0o600),
1013        )
1014        .await
1015        .map_err(|error| {
1016            format!(
1017                "secure isolated Codex config '{}': {error}",
1018                config_path.display()
1019            )
1020        })?;
1021        Ok(())
1022    }
1023
1024    async fn prepare_output_file(&self) -> Result<(), String> {
1025        let Some(path) = self.last_message_path() else {
1026            return Ok(());
1027        };
1028        if let Some(parent) = path.parent() {
1029            tokio::fs::create_dir_all(parent).await.map_err(|error| {
1030                format!("create Codex state dir '{}': {error}", parent.display())
1031            })?;
1032        }
1033        match tokio::fs::remove_file(&path).await {
1034            Ok(()) => Ok(()),
1035            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
1036            Err(error) => Err(format!(
1037                "remove stale Codex last-message file '{}': {error}",
1038                path.display()
1039            )),
1040        }
1041    }
1042
1043    fn handle_event(
1044        &self,
1045        value: Value,
1046        policy: &EffectiveCodexPolicy,
1047        events: &EventSink,
1048        state: &mut RunState,
1049    ) -> Option<String> {
1050        let event_type = value.get("type").and_then(Value::as_str).unwrap_or("");
1051        let mut captured_thread_id = None;
1052        match event_type {
1053            "thread.started" => {
1054                state.thread_id = value
1055                    .get("thread_id")
1056                    .and_then(Value::as_str)
1057                    .unwrap_or("")
1058                    .to_string();
1059                if !state.thread_id.is_empty() {
1060                    captured_thread_id = Some(state.thread_id.clone());
1061                }
1062                events.emit(json!({
1063                    "type": "runner_progress",
1064                    "session_id": state.thread_id,
1065                    "round_count": 0,
1066                    "executor": "codex",
1067                    "binary": self.binary,
1068                    "version": self.version,
1069                    "model": self.model,
1070                    "auth_mode": self.auth.mode().as_str(),
1071                    "codex_home_mode": self.codex_home_mode(),
1072                    "forward_env": self.forward_env,
1073                    "sandbox": policy.sandbox.as_str(),
1074                    "approval_policy": policy.approval_policy.as_str(),
1075                    "network_access": policy.network_access,
1076                    "policy_invocation": policy.invocation.as_str(),
1077                    "permission_profile": self.permissions.permission_profile.as_deref(),
1078                }));
1079            }
1080            "turn.started" => {
1081                state.turn_started = true;
1082                events.emit(event_json(AgentEvent::RunnerProgress {
1083                    session_id: state.session_id(),
1084                    round_count: 1,
1085                }));
1086            }
1087            "item.started" | "item.updated" | "item.completed" => {
1088                let phase = event_type.trim_start_matches("item.");
1089                if let Some(item) = value.get("item") {
1090                    handle_item(phase, item, events, state);
1091                    // Codex CLI 0.144 can omit the command_execution item when
1092                    // the OS sandbox itself rejects the command, even though
1093                    // it reports the exit status and errno to the model. Keep
1094                    // that safety failure visible instead of reducing the JSONL
1095                    // stream to an apparently successful agent message.
1096                    if phase == "completed"
1097                        && policy.sandbox != CodexSandbox::DangerFullAccess
1098                        && !state.tool_error_emitted
1099                        && item.get("type").and_then(Value::as_str) == Some("agent_message")
1100                        && item
1101                            .get("text")
1102                            .and_then(Value::as_str)
1103                            .is_some_and(looks_like_sandbox_denial)
1104                    {
1105                        let item_id = item
1106                            .get("id")
1107                            .and_then(Value::as_str)
1108                            .unwrap_or("codex-sandbox-denial");
1109                        let error = item.get("text").and_then(Value::as_str).unwrap_or(
1110                            "Codex reported that the OS sandbox denied a tool operation",
1111                        );
1112                        events.emit(event_json(AgentEvent::ToolError {
1113                            tool_call_id: format!("{item_id}-sandbox-denial"),
1114                            error: truncate_chars(error, TOOL_RESULT_TRUNCATE_CHARS),
1115                        }));
1116                        state.tool_error_emitted = true;
1117                    }
1118                }
1119            }
1120            "turn.completed" => {
1121                state.completed = true;
1122                state.usage = parse_usage(value.get("usage"));
1123                if let Some(text) = final_text_from_terminal(&value) {
1124                    state.last_agent_message = text;
1125                }
1126                events.emit(event_json(AgentEvent::Complete { usage: state.usage }));
1127            }
1128            "turn.failed" => {
1129                let message = error_message(&value, "Codex turn failed");
1130                state.failure = Some(message.clone());
1131                events.emit(event_json(AgentEvent::Error { message }));
1132            }
1133            "error" => {
1134                let message = error_message(&value, "Codex CLI error");
1135                state.failure = Some(message.clone());
1136                events.emit(event_json(AgentEvent::Error { message }));
1137            }
1138            other => {
1139                tracing::debug!(event_type = other, "codex: unrecognized JSONL event");
1140            }
1141        }
1142        captured_thread_id
1143    }
1144
1145    async fn read_last_message(&self) -> Option<String> {
1146        let path = self.last_message_path()?;
1147        let text = tokio::fs::read_to_string(path).await.ok()?;
1148        let trimmed = text.trim();
1149        (!trimmed.is_empty()).then(|| trimmed.to_string())
1150    }
1151
1152    async fn run_process(
1153        &self,
1154        prompt: &str,
1155        resume_id: Option<&str>,
1156        run_provider_token: Option<&str>,
1157        permission: bamboo_domain::PermissionModeResolution,
1158        events: &EventSink,
1159        cancel: &CancellationToken,
1160    ) -> (ChildOutcome, bool) {
1161        let policy = self
1162            .permissions
1163            .effective_for_session(permission, running_as_root());
1164        self.emit_policy_bootstrap(&policy, permission, events);
1165        // Warm workers reuse this executor across activations. Reassert the
1166        // isolated-home invariant at every run boundary so a prior Codex
1167        // process cannot leave an auth artifact or mutate provider routing for
1168        // the next child session.
1169        if let Err(error) = self.prepare_auth_home().await {
1170            return (ChildOutcome::error(error), false);
1171        }
1172        if let Err(error) = self.prepare_output_file().await {
1173            return (ChildOutcome::error(error), false);
1174        }
1175
1176        let mut child = match spawn_with_etxtbsy_retry(|| {
1177            self.build_command(run_provider_token, &policy, resume_id)
1178        })
1179        .await
1180        {
1181            Ok(child) => child,
1182            Err(error) => {
1183                return (
1184                    ChildOutcome::error(format!(
1185                        "spawn Codex CLI '{}': {error}; install with `npm i -g @openai/codex`, `brew install codex`, or an official GitHub release",
1186                        self.binary.display()
1187                    )),
1188                    false,
1189                );
1190            }
1191        };
1192        let Some(mut stdin) = child.stdin.take() else {
1193            terminate_child(&mut child).await;
1194            return (ChildOutcome::error("Codex child has no stdin pipe"), false);
1195        };
1196        let Some(stdout) = child.stdout.take() else {
1197            terminate_child(&mut child).await;
1198            return (ChildOutcome::error("Codex child has no stdout pipe"), false);
1199        };
1200        let stderr = child.stderr.take();
1201
1202        let write_result = tokio::select! {
1203            result = async {
1204                stdin.write_all(prompt.as_bytes()).await?;
1205                stdin.write_all(b"\n").await?;
1206                stdin.shutdown().await
1207            } => result,
1208            _ = cancel.cancelled() => {
1209                terminate_child(&mut child).await;
1210                events.emit(event_json(AgentEvent::Cancelled {
1211                    message: Some("Codex child cancelled".to_string()),
1212                }));
1213                return (ChildOutcome::cancelled(), false);
1214            }
1215        };
1216        if let Err(error) = write_result {
1217            terminate_child(&mut child).await;
1218            return (
1219                ChildOutcome::error(format!("write Codex prompt to stdin: {error}")),
1220                false,
1221            );
1222        }
1223        drop(stdin);
1224
1225        let stderr_tail = Arc::new(Mutex::new(String::new()));
1226        let stderr_task = stderr.map(|stderr| {
1227            let tail = stderr_tail.clone();
1228            tokio::spawn(async move { drain_stderr_tail(stderr, tail).await })
1229        });
1230
1231        let mut reader = BufReader::with_capacity(64 * 1024, stdout);
1232        let mut state = RunState::default();
1233        let mut read_error = None;
1234        let mut cancelled = false;
1235        loop {
1236            tokio::select! {
1237                _ = cancel.cancelled() => {
1238                    cancelled = true;
1239                    terminate_child(&mut child).await;
1240                    break;
1241                }
1242                line = read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES) => {
1243                    match line {
1244                        Ok(Some(bytes)) => {
1245                            if bytes.iter().all(u8::is_ascii_whitespace) {
1246                                continue;
1247                            }
1248                            match serde_json::from_slice::<Value>(&bytes) {
1249                                Ok(value) => {
1250                                    if let Some(thread_id) = self.handle_event(
1251                                        value,
1252                                        &policy,
1253                                        events,
1254                                        &mut state,
1255                                    ) {
1256                                        self.write_session_state(&thread_id).await;
1257                                    }
1258                                }
1259                                Err(error) => tracing::debug!(%error, "codex: skipping unparsable stdout line"),
1260                            }
1261                        }
1262                        Ok(None) => break,
1263                        Err(error) => {
1264                            read_error = Some(format!("Codex stdout read error: {error}"));
1265                            terminate_child(&mut child).await;
1266                            break;
1267                        }
1268                    }
1269                }
1270            }
1271        }
1272
1273        let status = if cancelled || read_error.is_some() {
1274            child.try_wait().ok().flatten()
1275        } else {
1276            match tokio::time::timeout(PROCESS_EXIT_WAIT, child.wait()).await {
1277                Ok(result) => result.ok(),
1278                Err(_) => {
1279                    terminate_child(&mut child).await;
1280                    child.try_wait().ok().flatten()
1281                }
1282            }
1283        };
1284        if let Some(task) = stderr_task {
1285            let _ = task.await;
1286        }
1287        let stderr = stderr_tail.lock().await.clone();
1288
1289        if cancelled {
1290            events.emit(event_json(AgentEvent::Cancelled {
1291                message: Some("Codex child cancelled".to_string()),
1292            }));
1293            return (ChildOutcome::cancelled(), false);
1294        }
1295
1296        let exited_without_turn = !state.turn_started && !state.completed;
1297        let outcome = if let Some(error) = read_error {
1298            ChildOutcome::error(error)
1299        } else if status.as_ref().is_some_and(|status| !status.success()) {
1300            ChildOutcome::error(format!(
1301                "Codex CLI exited with status {}; stderr tail: {}",
1302                status
1303                    .as_ref()
1304                    .map(ToString::to_string)
1305                    .unwrap_or_else(|| "<unknown>".to_string()),
1306                display_stderr_tail(&stderr)
1307            ))
1308        } else if let Some(error) = state.failure {
1309            ChildOutcome::error(format!(
1310                "{error}; stderr tail: {}",
1311                display_stderr_tail(&stderr)
1312            ))
1313        } else if !state.completed {
1314            ChildOutcome::error(format!(
1315                "Codex CLI exited without a turn.completed event; stderr tail: {}",
1316                display_stderr_tail(&stderr)
1317            ))
1318        } else {
1319            let final_text = if state.last_agent_message.trim().is_empty() {
1320                self.read_last_message().await
1321            } else {
1322                Some(state.last_agent_message)
1323            };
1324            match final_text {
1325                Some(text) => ChildOutcome::completed(text),
1326                None => ChildOutcome::error(
1327                    "Codex CLI completed without a final agent message or output-last-message file",
1328                ),
1329            }
1330        };
1331        (outcome, exited_without_turn)
1332    }
1333
1334    fn emit_policy_bootstrap(
1335        &self,
1336        policy: &EffectiveCodexPolicy,
1337        permission: bamboo_domain::PermissionModeResolution,
1338        events: &EventSink,
1339    ) {
1340        let executor_mapping = format!(
1341            "codex_exec:approval_policy={}",
1342            policy.approval_policy.as_str()
1343        );
1344        events.emit(json!({
1345            "type": "runner_progress",
1346            "session_id": "codex",
1347            "round_count": 0,
1348            "executor": "codex",
1349            "phase": "bootstrap",
1350            "binary": self.binary,
1351            "version": self.version,
1352            "model": self.model,
1353            "auth_mode": self.auth.mode().as_str(),
1354            "codex_home_mode": self.codex_home_mode(),
1355            "forward_env": self.forward_env,
1356            "sandbox": policy.sandbox.as_str(),
1357            "approval_policy": policy.approval_policy.as_str(),
1358            "network_access": policy.network_access,
1359            "policy_invocation": policy.invocation.as_str(),
1360            "permission_profile": self.permissions.permission_profile.as_deref(),
1361            "requested_mode": permission.requested.as_str(),
1362            "effective_mode": permission.effective.as_str(),
1363            "executor_mapping": executor_mapping,
1364        }));
1365        for warning in &policy.warnings {
1366            tracing::warn!(message = %warning, "Codex spawn policy warning");
1367            events.emit(json!({
1368                "type": "runner_progress",
1369                "session_id": "codex",
1370                "round_count": 0,
1371                "executor": "codex",
1372                "phase": "policy_warning",
1373                "level": "warning",
1374                "message": warning,
1375                "sandbox": policy.sandbox.as_str(),
1376                "approval_policy": policy.approval_policy.as_str(),
1377                "network_access": policy.network_access,
1378                "policy_invocation": policy.invocation.as_str(),
1379            }));
1380        }
1381    }
1382}
1383
1384#[async_trait]
1385impl ChildExecutor for CodexExecutor {
1386    async fn run(
1387        &self,
1388        spec: RunSpec,
1389        events: EventSink,
1390        mut steer: SteerInbox,
1391        cancel: CancellationToken,
1392    ) -> ChildOutcome {
1393        let logical_session_id = spec
1394            .permission_policy
1395            .as_ref()
1396            .map(|context| context.session_id.trim())
1397            .filter(|session_id| !session_id.is_empty())
1398            .or_else(|| {
1399                spec.logical_session
1400                    .as_ref()
1401                    .map(|identity| identity.session_id.trim())
1402                    .filter(|session_id| !session_id.is_empty())
1403            })
1404            .or_else(|| self.permissions.provisioned_session_id())
1405            .unwrap_or("codex")
1406            .to_string();
1407        let activation = match self
1408            .permissions
1409            .activation_permission(spec.permission_policy.as_ref())
1410        {
1411            Ok(activation) => activation,
1412            Err(error) => {
1413                events.emit(event_json(AgentEvent::Error {
1414                    message: error.clone(),
1415                }));
1416                return ChildOutcome::error(error);
1417            }
1418        };
1419        if let Some(reason) = activation.explicit_deny_reason.as_ref() {
1420            let message =
1421                format!("Codex exec cannot safely enforce Bamboo explicit-deny policy: {reason}");
1422            events.emit(event_json(AgentEvent::PermissionPostureActivated {
1423                session_id: logical_session_id,
1424                policy_revision: activation.policy_revision,
1425                requested_mode: activation.resolution.requested.as_str().to_string(),
1426                effective_mode: activation.resolution.effective.as_str().to_string(),
1427                executor_mapping: "codex_exec:blocked_explicit_deny".to_string(),
1428            }));
1429            events.emit(event_json(AgentEvent::Error {
1430                message: message.clone(),
1431            }));
1432            return ChildOutcome::error(message);
1433        }
1434        let bootstrap_policy = self
1435            .permissions
1436            .effective_for_session(activation.resolution, running_as_root());
1437        events.emit(event_json(AgentEvent::PermissionPostureActivated {
1438            session_id: logical_session_id,
1439            policy_revision: activation.policy_revision,
1440            requested_mode: activation.resolution.requested.as_str().to_string(),
1441            effective_mode: activation.resolution.effective.as_str().to_string(),
1442            executor_mapping: format!(
1443                "codex_exec:approval_policy={}",
1444                bootstrap_policy.approval_policy.as_str()
1445            ),
1446        }));
1447
1448        // `codex exec` v1 has no mid-turn steering channel. Drain the inbox so
1449        // transport senders cannot build an unbounded backlog.
1450        let steer_drain = tokio::spawn(async move { while steer.recv().await.is_some() {} });
1451        if spec.messages.is_empty() {
1452            self.delete_session_state().await;
1453        }
1454        let resume_id = if spec.messages.is_empty() {
1455            None
1456        } else {
1457            self.resolve_resume_id().await
1458        };
1459        let body = if resume_id.is_some() || spec.messages.is_empty() {
1460            spec.assignment.clone()
1461        } else {
1462            build_rehydrated_turn(&spec.messages, &spec.assignment)
1463        };
1464        let run_provider_token = spec
1465            .secrets
1466            .codex_provider_token
1467            .as_ref()
1468            .map(bamboo_subagent::proto::SecretValue::expose);
1469        let (outcome, exited_without_turn) = self
1470            .run_process(
1471                &body,
1472                resume_id.as_deref(),
1473                run_provider_token,
1474                activation.resolution,
1475                &events,
1476                &cancel,
1477            )
1478            .await;
1479        let outcome = if resume_id.is_some() && exited_without_turn {
1480            tracing::warn!(
1481                "codex: resume exited before turn progress; retrying once with rehydrated history"
1482            );
1483            events.emit(json!({
1484                "type": "runner_progress",
1485                "session_id": "codex",
1486                "round_count": 0,
1487                "executor": "codex",
1488                "phase": "resume_fallback",
1489                "message": "resume failed before turn progress; retrying once with rehydrated history",
1490            }));
1491            self.delete_session_state().await;
1492            let fallback_body = build_rehydrated_turn(&spec.messages, &spec.assignment);
1493            self.run_process(
1494                &fallback_body,
1495                None,
1496                run_provider_token,
1497                activation.resolution,
1498                &events,
1499                &cancel,
1500            )
1501            .await
1502            .0
1503        } else {
1504            outcome
1505        };
1506        steer_drain.abort();
1507        outcome
1508    }
1509}
1510
1511#[derive(Default)]
1512struct RunState {
1513    thread_id: String,
1514    turn_started: bool,
1515    completed: bool,
1516    failure: Option<String>,
1517    last_agent_message: String,
1518    usage: TokenUsage,
1519    tool_error_emitted: bool,
1520    started_items: HashSet<String>,
1521    item_text: HashMap<String, String>,
1522    item_output: HashMap<String, String>,
1523}
1524
1525impl RunState {
1526    fn session_id(&self) -> String {
1527        if self.thread_id.is_empty() {
1528            "codex".to_string()
1529        } else {
1530            self.thread_id.clone()
1531        }
1532    }
1533}
1534
1535fn handle_item(phase: &str, item: &Value, events: &EventSink, state: &mut RunState) {
1536    let item_type = item.get("type").and_then(Value::as_str).unwrap_or("");
1537    let item_id = item
1538        .get("id")
1539        .and_then(Value::as_str)
1540        .filter(|id| !id.is_empty())
1541        .map(str::to_string)
1542        .unwrap_or_else(|| format!("codex-{item_type}"));
1543
1544    match item_type {
1545        "agent_message" => {
1546            let text = item.get("text").and_then(Value::as_str).unwrap_or("");
1547            emit_text_delta(
1548                &item_id,
1549                text,
1550                &mut state.item_text,
1551                |delta| AgentEvent::Token { content: delta },
1552                events,
1553            );
1554            if phase == "completed" && !text.is_empty() {
1555                state.last_agent_message = text.to_string();
1556            }
1557        }
1558        "reasoning" => {
1559            let text = reasoning_text(item);
1560            emit_text_delta(
1561                &item_id,
1562                &text,
1563                &mut state.item_text,
1564                |delta| AgentEvent::ReasoningToken { content: delta },
1565                events,
1566            );
1567        }
1568        "command_execution" => {
1569            ensure_tool_started(
1570                &item_id,
1571                "Bash",
1572                json!({ "command": item.get("command").cloned().unwrap_or(Value::Null) }),
1573                events,
1574                state,
1575            );
1576            let output = item
1577                .get("aggregated_output")
1578                .and_then(Value::as_str)
1579                .unwrap_or("");
1580            emit_tool_output_delta(&item_id, output, events, state);
1581            if phase == "completed" {
1582                let exit_code = item.get("exit_code").and_then(Value::as_i64);
1583                let successful = exit_code == Some(0)
1584                    && item.get("status").and_then(Value::as_str) != Some("failed");
1585                if successful {
1586                    events.emit(event_json(AgentEvent::ToolComplete {
1587                        tool_call_id: item_id,
1588                        result: ToolResult::text(
1589                            true,
1590                            truncate_chars(output, TOOL_RESULT_TRUNCATE_CHARS),
1591                        ),
1592                    }));
1593                } else {
1594                    let error = if output.trim().is_empty() {
1595                        format!("command failed with exit code {exit_code:?}")
1596                    } else {
1597                        truncate_chars(output, TOOL_RESULT_TRUNCATE_CHARS)
1598                    };
1599                    events.emit(event_json(AgentEvent::ToolError {
1600                        tool_call_id: item_id,
1601                        error,
1602                    }));
1603                    state.tool_error_emitted = true;
1604                }
1605            }
1606        }
1607        "file_change" => {
1608            let detail = item
1609                .get("changes")
1610                .or_else(|| item.get("patch"))
1611                .cloned()
1612                .unwrap_or_else(|| item.clone());
1613            ensure_tool_started(
1614                &item_id,
1615                "ApplyPatch",
1616                json!({ "changes": detail }),
1617                events,
1618                state,
1619            );
1620            if phase == "completed" {
1621                complete_structured_tool(&item_id, item, events, state);
1622            }
1623        }
1624        "mcp_tool_call" => {
1625            let server = item.get("server").and_then(Value::as_str).unwrap_or("mcp");
1626            let tool = item
1627                .get("tool")
1628                .or_else(|| item.get("name"))
1629                .and_then(Value::as_str)
1630                .unwrap_or("tool");
1631            let tool_name = format!("{server}::{tool}");
1632            let arguments = item
1633                .get("arguments")
1634                .or_else(|| item.get("input"))
1635                .cloned()
1636                .unwrap_or_else(|| json!({}));
1637            ensure_tool_started(&item_id, &tool_name, arguments, events, state);
1638            if phase == "completed" {
1639                complete_structured_tool(&item_id, item, events, state);
1640            }
1641        }
1642        "web_search" => {
1643            let query = item.get("query").cloned().unwrap_or(Value::Null);
1644            ensure_tool_started(
1645                &item_id,
1646                "WebSearch",
1647                json!({ "query": query }),
1648                events,
1649                state,
1650            );
1651            if phase == "completed" {
1652                complete_structured_tool(&item_id, item, events, state);
1653            }
1654        }
1655        "todo_list" => {
1656            events.emit(json!({
1657                "type": "runner_progress",
1658                "session_id": state.session_id(),
1659                "round_count": 1,
1660                "codex_item_type": "todo_list",
1661                "item": item,
1662            }));
1663        }
1664        other => {
1665            tracing::debug!(item_type = other, phase, "codex: unrecognized item type");
1666        }
1667    }
1668}
1669
1670fn ensure_tool_started(
1671    item_id: &str,
1672    tool_name: &str,
1673    arguments: Value,
1674    events: &EventSink,
1675    state: &mut RunState,
1676) {
1677    if state.started_items.insert(item_id.to_string()) {
1678        events.emit(event_json(AgentEvent::ToolStart {
1679            tool_call_id: item_id.to_string(),
1680            tool_name: tool_name.to_string(),
1681            arguments,
1682        }));
1683    }
1684}
1685
1686fn emit_tool_output_delta(item_id: &str, output: &str, events: &EventSink, state: &mut RunState) {
1687    let previous = state.item_output.entry(item_id.to_string()).or_default();
1688    let delta = if output.starts_with(previous.as_str()) {
1689        &output[previous.len()..]
1690    } else {
1691        output
1692    };
1693    if !delta.is_empty() {
1694        events.emit(event_json(AgentEvent::ToolToken {
1695            tool_call_id: item_id.to_string(),
1696            content: delta.to_string(),
1697        }));
1698    }
1699    *previous = output.to_string();
1700}
1701
1702fn complete_structured_tool(item_id: &str, item: &Value, events: &EventSink, state: &mut RunState) {
1703    if let Some(error) = item.get("error").filter(|value| !value.is_null()) {
1704        events.emit(event_json(AgentEvent::ToolError {
1705            tool_call_id: item_id.to_string(),
1706            error: value_text(error),
1707        }));
1708        state.tool_error_emitted = true;
1709        return;
1710    }
1711    let result = item
1712        .get("result")
1713        .or_else(|| item.get("output"))
1714        .cloned()
1715        .unwrap_or_else(|| item.clone());
1716    events.emit(event_json(AgentEvent::ToolComplete {
1717        tool_call_id: item_id.to_string(),
1718        result: ToolResult::text(
1719            true,
1720            truncate_chars(&value_text(&result), TOOL_RESULT_TRUNCATE_CHARS),
1721        ),
1722    }));
1723}
1724
1725fn looks_like_sandbox_denial(text: &str) -> bool {
1726    let normalized = text.to_ascii_lowercase();
1727    let denied = normalized.contains("operation not permitted")
1728        || normalized.contains("permission denied")
1729        || normalized.contains("read-only file system");
1730    let operation = normalized.contains("sandbox")
1731        || normalized.contains("command")
1732        || normalized.contains("write")
1733        || normalized.contains("exit status");
1734    denied && operation
1735}
1736
1737fn emit_text_delta<F>(
1738    item_id: &str,
1739    text: &str,
1740    seen: &mut HashMap<String, String>,
1741    build: F,
1742    events: &EventSink,
1743) where
1744    F: FnOnce(String) -> AgentEvent,
1745{
1746    let previous = seen.entry(item_id.to_string()).or_default();
1747    let delta = if text.starts_with(previous.as_str()) {
1748        &text[previous.len()..]
1749    } else {
1750        text
1751    };
1752    if !delta.is_empty() {
1753        events.emit(event_json(build(delta.to_string())));
1754    }
1755    *previous = text.to_string();
1756}
1757
1758fn reasoning_text(item: &Value) -> String {
1759    match item.get("text").or_else(|| item.get("summary")) {
1760        Some(Value::String(text)) => text.clone(),
1761        Some(Value::Array(parts)) => parts
1762            .iter()
1763            .filter_map(|part| {
1764                part.as_str()
1765                    .or_else(|| part.get("text").and_then(Value::as_str))
1766            })
1767            .collect::<Vec<_>>()
1768            .join("\n"),
1769        Some(value) => value.to_string(),
1770        None => String::new(),
1771    }
1772}
1773
1774fn parse_usage(value: Option<&Value>) -> TokenUsage {
1775    let input = value
1776        .and_then(|usage| usage.get("input_tokens"))
1777        .and_then(Value::as_u64)
1778        .unwrap_or(0);
1779    let output = value
1780        .and_then(|usage| usage.get("output_tokens"))
1781        .and_then(Value::as_u64)
1782        .unwrap_or(0);
1783    TokenUsage {
1784        prompt_tokens: input,
1785        completion_tokens: output,
1786        total_tokens: input.saturating_add(output),
1787    }
1788}
1789
1790fn final_text_from_terminal(value: &Value) -> Option<String> {
1791    ["final_output", "output_text", "result"]
1792        .into_iter()
1793        .find_map(|key| value.get(key).and_then(Value::as_str))
1794        .filter(|text| !text.is_empty())
1795        .map(str::to_string)
1796}
1797
1798fn error_message(value: &Value, fallback: &str) -> String {
1799    value
1800        .pointer("/error/message")
1801        .or_else(|| value.get("message"))
1802        .or_else(|| value.get("error"))
1803        .map(value_text)
1804        .filter(|message| !message.is_empty())
1805        .unwrap_or_else(|| fallback.to_string())
1806}
1807
1808fn value_text(value: &Value) -> String {
1809    value
1810        .as_str()
1811        .map(str::to_string)
1812        .unwrap_or_else(|| value.to_string())
1813}
1814
1815fn display_stderr_tail(stderr: &str) -> &str {
1816    let trimmed = stderr.trim();
1817    if trimmed.is_empty() {
1818        "<empty>"
1819    } else {
1820        trimmed
1821    }
1822}
1823
1824fn event_json(event: AgentEvent) -> Value {
1825    serde_json::to_value(event).unwrap_or_else(|_| json!({}))
1826}
1827
1828fn truncate_chars(text: &str, max_chars: usize) -> String {
1829    if text.chars().count() <= max_chars {
1830        return text.to_string();
1831    }
1832    let head: String = text.chars().take(max_chars).collect();
1833    let dropped = text.chars().count().saturating_sub(max_chars);
1834    format!("{head}\n… [truncated, {dropped} more chars]")
1835}
1836
1837fn has_git_metadata(path: &Path) -> bool {
1838    path.ancestors()
1839        .any(|ancestor| ancestor.join(".git").exists())
1840}
1841
1842fn validate_codex_base_url(mode: CodexAuthMode, raw: Option<&str>) -> Result<(), String> {
1843    match mode {
1844        CodexAuthMode::Custom | CodexAuthMode::Bamboo => {
1845            let raw = raw
1846                .map(str::trim)
1847                .filter(|value| !value.is_empty())
1848                .ok_or_else(|| format!("Codex auth mode '{mode:?}' requires codex_base_url"))?;
1849            let parsed = url::Url::parse(raw)
1850                .map_err(|error| format!("invalid Codex base URL '{raw}': {error}"))?;
1851            if !matches!(parsed.scheme(), "http" | "https") || parsed.host_str().is_none() {
1852                return Err("Codex base URL must be an absolute HTTP(S) URL".to_string());
1853            }
1854            if !parsed.username().is_empty()
1855                || parsed.password().is_some()
1856                || parsed.query().is_some()
1857                || parsed.fragment().is_some()
1858            {
1859                return Err(
1860                    "Codex base URL must not contain credentials, query parameters, or a fragment"
1861                        .to_string(),
1862                );
1863            }
1864        }
1865        CodexAuthMode::Inherit | CodexAuthMode::ApiKey => {
1866            if raw.is_some_and(|value| !value.trim().is_empty()) {
1867                return Err(
1868                    "codex_base_url is only valid with custom or bamboo auth mode".to_string(),
1869                );
1870            }
1871        }
1872    }
1873    Ok(())
1874}
1875
1876fn validate_codex_forward_env(mode: CodexAuthMode, names: &[String]) -> Result<(), String> {
1877    let mut seen = HashSet::new();
1878    for name in names {
1879        let valid = name
1880            .chars()
1881            .next()
1882            .is_some_and(|first| first == '_' || first.is_ascii_alphabetic())
1883            && name
1884                .chars()
1885                .all(|character| character == '_' || character.is_ascii_alphanumeric());
1886        if !valid {
1887            return Err(format!("invalid Codex forward_env name '{name}'"));
1888        }
1889        if !seen.insert(name.as_str()) {
1890            return Err(format!("duplicate Codex forward_env name '{name}'"));
1891        }
1892        if name.starts_with("CODEX_") || name == CODEX_PROVIDER_ENV {
1893            return Err(format!(
1894                "Codex forward_env may not override reserved variable '{name}'"
1895            ));
1896        }
1897    }
1898    let forwards_openai = seen.contains("OPENAI_API_KEY");
1899    if mode == CodexAuthMode::ApiKey && !forwards_openai {
1900        return Err(
1901            "Codex api_key auth requires explicit OPENAI_API_KEY in codex_forward_env".to_string(),
1902        );
1903    }
1904    if mode != CodexAuthMode::ApiKey && forwards_openai {
1905        return Err("OPENAI_API_KEY may only be forwarded in Codex api_key auth mode".to_string());
1906    }
1907    Ok(())
1908}
1909
1910async fn spawn_with_etxtbsy_retry(
1911    mut build: impl FnMut() -> Result<Command, String>,
1912) -> Result<Child, String> {
1913    let mut last_error = None;
1914    for _ in 0..5 {
1915        match build()?.spawn() {
1916            Ok(child) => return Ok(child),
1917            Err(error) if error.raw_os_error() == Some(26) => {
1918                last_error = Some(error);
1919                tokio::time::sleep(Duration::from_millis(10)).await;
1920            }
1921            Err(error) => return Err(error.to_string()),
1922        }
1923    }
1924    Err(last_error
1925        .expect("retry loop records ETXTBSY before exhausting")
1926        .to_string())
1927}
1928
1929#[derive(Clone, Copy)]
1930enum ProcessSignal {
1931    Term,
1932    Kill,
1933}
1934
1935#[cfg(unix)]
1936fn signal_process_group(pgid: libc::pid_t, signal: ProcessSignal) {
1937    let signal = match signal {
1938        ProcessSignal::Term => libc::SIGTERM,
1939        ProcessSignal::Kill => libc::SIGKILL,
1940    };
1941    // SAFETY: the pgid is from our child and build_command created a new
1942    // process group with that child as leader. ESRCH is harmless.
1943    unsafe {
1944        libc::kill(-pgid, signal);
1945    }
1946}
1947
1948#[cfg(unix)]
1949fn signal_process(pid: libc::pid_t, signal: ProcessSignal) {
1950    let signal = match signal {
1951        ProcessSignal::Term => libc::SIGTERM,
1952        ProcessSignal::Kill => libc::SIGKILL,
1953    };
1954    // SAFETY: the pid comes from the OS process table. ESRCH is harmless.
1955    unsafe {
1956        libc::kill(pid, signal);
1957    }
1958}
1959
1960#[cfg(unix)]
1961fn descendant_processes(root: libc::pid_t) -> Vec<libc::pid_t> {
1962    use sysinfo::{ProcessesToUpdate, System};
1963
1964    let mut system = System::new();
1965    system.refresh_processes(ProcessesToUpdate::All, true);
1966    let mut known = HashSet::from([root]);
1967    let mut descendants = Vec::new();
1968    loop {
1969        let mut added = false;
1970        for (pid, process) in system.processes() {
1971            let pid = pid.as_u32() as libc::pid_t;
1972            let parent = process
1973                .parent()
1974                .map(|parent| parent.as_u32() as libc::pid_t);
1975            if !known.contains(&pid) && parent.is_some_and(|parent| known.contains(&parent)) {
1976                known.insert(pid);
1977                descendants.push(pid);
1978                added = true;
1979            }
1980        }
1981        if !added {
1982            break;
1983        }
1984    }
1985    descendants
1986}
1987
1988#[cfg(unix)]
1989fn process_exists(pid: libc::pid_t) -> bool {
1990    // SAFETY: signal 0 only probes whether the process exists.
1991    if unsafe { libc::kill(pid, 0) } == 0 {
1992        return true;
1993    }
1994    std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
1995}
1996
1997#[cfg(unix)]
1998fn process_group_exists(pgid: libc::pid_t) -> bool {
1999    // SAFETY: signal 0 only probes whether the process group exists.
2000    if unsafe { libc::kill(-pgid, 0) } == 0 {
2001        return true;
2002    }
2003    std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
2004}
2005
2006#[cfg(unix)]
2007pub(crate) async fn terminate_child(child: &mut Child) {
2008    let Some(pgid) = child.id().map(|pid| pid as libc::pid_t) else {
2009        return;
2010    };
2011    // Codex may launch a sandbox/tool command in a NEW process group. Capture
2012    // the full descendant tree before terminating the Codex leader; otherwise
2013    // those commands are reparented to init and escape a group-only signal.
2014    let descendants = descendant_processes(pgid);
2015    for &pid in descendants.iter().rev() {
2016        signal_process(pid, ProcessSignal::Term);
2017    }
2018    signal_process_group(pgid, ProcessSignal::Term);
2019    let deadline = tokio::time::Instant::now() + SIGTERM_WAIT;
2020    loop {
2021        // Reap the leader when it exits, but keep tracking its process group
2022        // and any descendants Codex placed in independent groups.
2023        let _ = child.try_wait();
2024        if !process_group_exists(pgid) && !descendants.iter().copied().any(process_exists) {
2025            let _ = child.wait().await;
2026            return;
2027        }
2028        if tokio::time::Instant::now() >= deadline {
2029            break;
2030        }
2031        tokio::time::sleep(Duration::from_millis(25)).await;
2032    }
2033    for &pid in descendants.iter().rev() {
2034        signal_process(pid, ProcessSignal::Kill);
2035    }
2036    signal_process_group(pgid, ProcessSignal::Kill);
2037    let _ = child.start_kill();
2038    let _ = child.wait().await;
2039}
2040
2041#[cfg(not(unix))]
2042pub(crate) async fn terminate_child(child: &mut Child) {
2043    if tokio::time::timeout(SIGTERM_WAIT, child.wait())
2044        .await
2045        .is_ok()
2046    {
2047        return;
2048    }
2049    let _ = child.start_kill();
2050    let _ = child.wait().await;
2051}
2052
2053pub(crate) async fn read_bounded_line<R>(
2054    reader: &mut R,
2055    max_bytes: usize,
2056) -> std::io::Result<Option<Vec<u8>>>
2057where
2058    R: tokio::io::AsyncBufRead + Unpin,
2059{
2060    let mut output = Vec::new();
2061    loop {
2062        let (found_newline, consumed) = {
2063            let available = reader.fill_buf().await?;
2064            if available.is_empty() {
2065                return Ok(if output.is_empty() {
2066                    None
2067                } else {
2068                    Some(output)
2069                });
2070            }
2071            match available.iter().position(|byte| *byte == b'\n') {
2072                Some(position) => {
2073                    output.extend_from_slice(&available[..position]);
2074                    (true, position + 1)
2075                }
2076                None => {
2077                    output.extend_from_slice(available);
2078                    (false, available.len())
2079                }
2080            }
2081        };
2082        reader.consume(consumed);
2083        if output.len() > max_bytes {
2084            return Err(std::io::Error::new(
2085                std::io::ErrorKind::InvalidData,
2086                format!("stdout line exceeded {max_bytes} bytes"),
2087            ));
2088        }
2089        if found_newline {
2090            return Ok(Some(output));
2091        }
2092    }
2093}
2094
2095async fn drain_stderr_tail(stderr: tokio::process::ChildStderr, tail: Arc<Mutex<String>>) {
2096    let mut reader = BufReader::new(stderr);
2097    let mut buffer = Vec::new();
2098    loop {
2099        buffer.clear();
2100        match reader.read_until(b'\n', &mut buffer).await {
2101            Ok(0) | Err(_) => return,
2102            Ok(_) => {
2103                let mut tail = tail.lock().await;
2104                tail.push_str(&String::from_utf8_lossy(&buffer));
2105                if tail.len() > STDERR_TAIL_BYTES {
2106                    let excess = tail.len() - STDERR_TAIL_BYTES;
2107                    let cut = tail
2108                        .char_indices()
2109                        .map(|(index, _)| index)
2110                        .find(|index| *index >= excess)
2111                        .unwrap_or(tail.len());
2112                    tail.drain(..cut);
2113                }
2114            }
2115        }
2116    }
2117}
2118
2119#[cfg(test)]
2120mod tests {
2121    use super::*;
2122    use bamboo_subagent::proto::TerminalStatus;
2123    use bamboo_subagent::provision::ScopedCredential;
2124
2125    fn resolved_auth(
2126        mode: Option<&str>,
2127        base_url: Option<&str>,
2128        provider_key_ref: Option<&str>,
2129        credentials: &[ScopedCredential],
2130        forward_env: &[&str],
2131    ) -> Result<CodexAuthConfig, String> {
2132        resolve_codex_auth_config(
2133            mode,
2134            false,
2135            base_url.map(str::to_string),
2136            Some("responses".to_string()),
2137            provider_key_ref,
2138            credentials,
2139            &forward_env
2140                .iter()
2141                .map(|name| name.to_string())
2142                .collect::<Vec<_>>(),
2143        )
2144    }
2145
2146    fn command_env(command: &Command, name: &str) -> Option<String> {
2147        command
2148            .as_std()
2149            .get_envs()
2150            .find(|(key, _)| *key == name)
2151            .and_then(|(_, value)| value)
2152            .map(|value| value.to_string_lossy().into_owned())
2153    }
2154
2155    fn command_args(command: &Command) -> Vec<String> {
2156        command
2157            .as_std()
2158            .get_args()
2159            .map(|argument| argument.to_string_lossy().into_owned())
2160            .collect()
2161    }
2162
2163    fn auth_error(result: Result<CodexAuthConfig, String>) -> String {
2164        match result {
2165            Err(error) => error,
2166            Ok(_) => panic!("expected Codex auth configuration to be rejected"),
2167        }
2168    }
2169
2170    fn permissions(
2171        sandbox: Option<&str>,
2172        approval_policy: Option<&str>,
2173        network_access: bool,
2174        allow_danger_bypass: bool,
2175        permission_profile: Option<&str>,
2176        provisioned_bypass: bool,
2177        workspace_owned: bool,
2178    ) -> CodexPermissionConfig {
2179        resolve_codex_permission_config(
2180            sandbox,
2181            approval_policy,
2182            network_access,
2183            allow_danger_bypass,
2184            permission_profile.map(str::to_string),
2185            provisioned_bypass,
2186            workspace_owned,
2187        )
2188        .unwrap()
2189    }
2190
2191    fn default_permissions() -> CodexPermissionConfig {
2192        permissions(None, None, false, false, None, false, false)
2193    }
2194
2195    fn default_policy(executor: &CodexExecutor) -> EffectiveCodexPolicy {
2196        executor.permissions.effective(false, false)
2197    }
2198
2199    #[test]
2200    fn run_permission_context_replaces_provisioned_auto_instead_of_sticking() {
2201        let permissions = default_permissions().with_provisioned_permission_resolution(
2202            bamboo_domain::resolve_permission_mode(
2203                bamboo_domain::SessionPermissionMode::Auto,
2204                bamboo_domain::PermissionMode::Default,
2205            ),
2206            "provisioned-codex".to_string(),
2207        );
2208        assert_eq!(
2209            permissions.activation_permission_resolution(None).unwrap(),
2210            bamboo_domain::resolve_permission_mode(
2211                bamboo_domain::SessionPermissionMode::Auto,
2212                bamboo_domain::PermissionMode::Default,
2213            )
2214        );
2215
2216        let explicit_default = bamboo_subagent::proto::PermissionPolicyContext {
2217            revision: 3,
2218            requested_mode: "default".to_string(),
2219            effective_mode: "default".to_string(),
2220            bypass_permissions: false,
2221            auto_approve_permissions: false,
2222            session_id: "explicit-default".to_string(),
2223            workspace_path: None,
2224            inherit_session_grants: false,
2225            policy: json!({}),
2226        };
2227        assert_eq!(
2228            permissions
2229                .activation_permission_resolution(Some(&explicit_default))
2230                .unwrap(),
2231            bamboo_domain::resolve_permission_mode(
2232                bamboo_domain::SessionPermissionMode::Default,
2233                bamboo_domain::PermissionMode::Default,
2234            )
2235        );
2236
2237        let explicit_bypass = bamboo_subagent::proto::PermissionPolicyContext {
2238            requested_mode: "bypass".to_string(),
2239            effective_mode: "bypass".to_string(),
2240            bypass_permissions: true,
2241            ..explicit_default
2242        };
2243        assert_eq!(
2244            permissions
2245                .activation_permission_resolution(Some(&explicit_bypass))
2246                .unwrap(),
2247            bamboo_domain::resolve_permission_mode(
2248                bamboo_domain::SessionPermissionMode::Bypass,
2249                bamboo_domain::PermissionMode::Default,
2250            )
2251        );
2252    }
2253
2254    fn fixture_executor() -> CodexExecutor {
2255        CodexExecutor {
2256            binary: PathBuf::from("/usr/local/bin/codex"),
2257            version: "codex-cli 0.144.5".to_string(),
2258            model: None,
2259            permissions: default_permissions(),
2260            workspace: None,
2261            state_dir: None,
2262            forward_env: Vec::new(),
2263            auth: CodexAuthConfig::inherit(),
2264        }
2265    }
2266
2267    fn map_fixture(input: &str) -> (RunState, Vec<Value>) {
2268        let executor = fixture_executor();
2269        let (sink, mut rx) = EventSink::channel();
2270        let mut state = RunState::default();
2271        let policy = default_policy(&executor);
2272        for line in input.lines().filter(|line| !line.trim().is_empty()) {
2273            executor.handle_event(
2274                serde_json::from_str(line).unwrap(),
2275                &policy,
2276                &sink,
2277                &mut state,
2278            );
2279        }
2280        let mut events = Vec::new();
2281        while let Ok(event) = rx.try_recv() {
2282            events.push(event);
2283        }
2284        (state, events)
2285    }
2286
2287    #[test]
2288    fn permission_profile_mapping_table_is_explicit_and_double_gated() {
2289        struct Case {
2290            name: &'static str,
2291            permissions: CodexPermissionConfig,
2292            parent_bypass: bool,
2293            is_root: bool,
2294            sandbox: CodexSandbox,
2295            approval: CodexApprovalPolicy,
2296            invocation: CodexPolicyInvocation,
2297            network: bool,
2298            warning: bool,
2299        }
2300
2301        let cases = [
2302            Case {
2303                name: "default",
2304                permissions: permissions(None, None, false, false, None, false, false),
2305                parent_bypass: false,
2306                is_root: false,
2307                sandbox: CodexSandbox::WorkspaceWrite,
2308                approval: CodexApprovalPolicy::Never,
2309                invocation: CodexPolicyInvocation::Explicit,
2310                network: false,
2311                warning: false,
2312            },
2313            Case {
2314                name: "restricted",
2315                permissions: permissions(
2316                    None,
2317                    None,
2318                    false,
2319                    false,
2320                    Some("restricted"),
2321                    false,
2322                    false,
2323                ),
2324                parent_bypass: false,
2325                is_root: false,
2326                sandbox: CodexSandbox::WorkspaceWrite,
2327                approval: CodexApprovalPolicy::Never,
2328                invocation: CodexPolicyInvocation::Explicit,
2329                network: false,
2330                warning: false,
2331            },
2332            Case {
2333                name: "research",
2334                permissions: permissions(None, None, false, false, Some("research"), false, false),
2335                parent_bypass: false,
2336                is_root: false,
2337                sandbox: CodexSandbox::ReadOnly,
2338                approval: CodexApprovalPolicy::Never,
2339                invocation: CodexPolicyInvocation::Explicit,
2340                network: false,
2341                warning: false,
2342            },
2343            Case {
2344                name: "network workspace",
2345                permissions: permissions(None, None, true, false, None, false, false),
2346                parent_bypass: false,
2347                is_root: false,
2348                sandbox: CodexSandbox::WorkspaceWrite,
2349                approval: CodexApprovalPolicy::Never,
2350                invocation: CodexPolicyInvocation::Explicit,
2351                network: true,
2352                warning: false,
2353            },
2354            Case {
2355                name: "ordinary bypass stays sandboxed",
2356                permissions: permissions(None, None, false, false, None, false, false),
2357                parent_bypass: true,
2358                is_root: false,
2359                sandbox: CodexSandbox::WorkspaceWrite,
2360                approval: CodexApprovalPolicy::Never,
2361                invocation: CodexPolicyInvocation::FullAuto,
2362                network: false,
2363                warning: false,
2364            },
2365            Case {
2366                name: "config-only danger opt-in downgrades",
2367                permissions: permissions(None, None, false, true, None, false, false),
2368                parent_bypass: false,
2369                is_root: false,
2370                sandbox: CodexSandbox::WorkspaceWrite,
2371                approval: CodexApprovalPolicy::Never,
2372                invocation: CodexPolicyInvocation::FullAuto,
2373                network: false,
2374                warning: true,
2375            },
2376            Case {
2377                name: "both danger gates",
2378                permissions: permissions(None, None, false, true, None, false, false),
2379                parent_bypass: true,
2380                is_root: false,
2381                sandbox: CodexSandbox::DangerFullAccess,
2382                approval: CodexApprovalPolicy::Never,
2383                invocation: CodexPolicyInvocation::DangerBypass,
2384                network: true,
2385                warning: true,
2386            },
2387            Case {
2388                name: "root danger request downgrades",
2389                permissions: permissions(
2390                    Some("danger-full-access"),
2391                    None,
2392                    false,
2393                    true,
2394                    None,
2395                    false,
2396                    false,
2397                ),
2398                parent_bypass: true,
2399                is_root: true,
2400                sandbox: CodexSandbox::WorkspaceWrite,
2401                approval: CodexApprovalPolicy::Never,
2402                invocation: CodexPolicyInvocation::FullAuto,
2403                network: false,
2404                warning: true,
2405            },
2406            Case {
2407                name: "explicit safe override",
2408                permissions: permissions(
2409                    Some("workspace-write"),
2410                    Some("on-failure"),
2411                    false,
2412                    false,
2413                    Some("read-only"),
2414                    false,
2415                    false,
2416                ),
2417                parent_bypass: false,
2418                is_root: false,
2419                sandbox: CodexSandbox::WorkspaceWrite,
2420                approval: CodexApprovalPolicy::OnFailure,
2421                invocation: CodexPolicyInvocation::Explicit,
2422                network: false,
2423                warning: false,
2424            },
2425        ];
2426
2427        for case in cases {
2428            let actual = case.permissions.effective(case.parent_bypass, case.is_root);
2429            assert_eq!(actual.sandbox, case.sandbox, "{} sandbox", case.name);
2430            assert_eq!(
2431                actual.approval_policy, case.approval,
2432                "{} approval",
2433                case.name
2434            );
2435            assert_eq!(
2436                actual.invocation, case.invocation,
2437                "{} invocation",
2438                case.name
2439            );
2440            assert_eq!(actual.network_access, case.network, "{} network", case.name);
2441            assert_eq!(
2442                !actual.warnings.is_empty(),
2443                case.warning,
2444                "{} warning",
2445                case.name
2446            );
2447        }
2448    }
2449
2450    #[test]
2451    fn permission_configuration_rejects_interactive_and_incoherent_values() {
2452        assert!(resolve_codex_permission_config(
2453            None,
2454            Some("on-request"),
2455            false,
2456            false,
2457            None,
2458            false,
2459            false,
2460        )
2461        .unwrap_err()
2462        .contains("non-interactive"));
2463        assert!(resolve_codex_permission_config(
2464            Some("read-only"),
2465            None,
2466            true,
2467            false,
2468            None,
2469            false,
2470            false,
2471        )
2472        .unwrap_err()
2473        .contains("effective workspace-write"));
2474        assert!(resolve_codex_permission_config(
2475            None,
2476            None,
2477            true,
2478            false,
2479            Some("research".to_string()),
2480            false,
2481            false,
2482        )
2483        .unwrap_err()
2484        .contains("effective workspace-write"));
2485    }
2486
2487    #[test]
2488    fn auto_forces_never_without_widening_the_codex_sandbox() {
2489        let permissions = permissions(
2490            Some("read-only"),
2491            Some("on-failure"),
2492            false,
2493            true,
2494            Some("guardian"),
2495            false,
2496            true,
2497        );
2498
2499        let baseline = permissions.effective(false, false);
2500        assert_eq!(baseline.sandbox, CodexSandbox::ReadOnly);
2501        assert_eq!(baseline.approval_policy, CodexApprovalPolicy::OnFailure);
2502
2503        let auto = permissions.effective_for_session(
2504            bamboo_domain::resolve_permission_mode(
2505                bamboo_domain::SessionPermissionMode::Auto,
2506                bamboo_domain::PermissionMode::Default,
2507            ),
2508            false,
2509        );
2510        assert_eq!(auto.sandbox, CodexSandbox::ReadOnly);
2511        assert_eq!(auto.approval_policy, CodexApprovalPolicy::Never);
2512        assert_eq!(auto.invocation, CodexPolicyInvocation::Explicit);
2513        assert!(!auto.network_access);
2514    }
2515
2516    #[test]
2517    fn command_flags_match_effective_policy_and_workspace_ownership() {
2518        let workspace = tempfile::tempdir().unwrap();
2519        let mut executor = fixture_executor();
2520        executor.workspace = Some(workspace.path().to_string_lossy().into_owned());
2521
2522        let policy = executor.permissions.effective(false, false);
2523        let args = command_args(&executor.build_command(None, &policy, None).unwrap());
2524        assert!(args
2525            .windows(2)
2526            .any(|pair| pair == ["--sandbox", "workspace-write"]));
2527        assert!(args
2528            .windows(2)
2529            .any(|pair| pair == ["--config", "approval_policy=\"never\""]));
2530        assert!(!args
2531            .iter()
2532            .any(|argument| argument == "--skip-git-repo-check"));
2533
2534        executor.permissions = permissions(None, None, false, false, None, false, true);
2535        let policy = executor.permissions.effective(false, false);
2536        let args = command_args(&executor.build_command(None, &policy, None).unwrap());
2537        assert!(args
2538            .iter()
2539            .any(|argument| argument == "--skip-git-repo-check"));
2540
2541        executor.permissions = permissions(None, None, true, false, None, false, true);
2542        let policy = executor.permissions.effective(false, false);
2543        let args = command_args(&executor.build_command(None, &policy, None).unwrap());
2544        assert!(args
2545            .windows(2)
2546            .any(|pair| { pair == ["--config", "sandbox_workspace_write.network_access=true"] }));
2547
2548        executor.permissions = permissions(None, None, false, false, None, false, true);
2549        let policy = executor.permissions.effective(true, false);
2550        let args = command_args(&executor.build_command(None, &policy, None).unwrap());
2551        assert!(args.iter().any(|argument| argument == "--full-auto"));
2552        assert!(!args
2553            .iter()
2554            .any(|argument| argument == "--dangerously-bypass-approvals-and-sandbox"));
2555
2556        executor.permissions = permissions(
2557            Some("danger-full-access"),
2558            None,
2559            false,
2560            true,
2561            None,
2562            false,
2563            true,
2564        );
2565        let policy = executor.permissions.effective(true, false);
2566        let args = command_args(&executor.build_command(None, &policy, None).unwrap());
2567        assert!(args
2568            .iter()
2569            .any(|argument| argument == "--dangerously-bypass-approvals-and-sandbox"));
2570        assert!(!args.iter().any(|argument| argument == "--full-auto"));
2571    }
2572
2573    #[test]
2574    fn danger_bypass_emits_loud_audit_warning() {
2575        let mut executor = fixture_executor();
2576        executor.permissions = permissions(
2577            Some("danger-full-access"),
2578            None,
2579            false,
2580            true,
2581            Some("bypass"),
2582            false,
2583            false,
2584        );
2585        let policy = executor.permissions.effective(true, false);
2586        let (sink, mut receiver) = EventSink::channel();
2587        executor.emit_policy_bootstrap(
2588            &policy,
2589            bamboo_domain::resolve_permission_mode(
2590                bamboo_domain::SessionPermissionMode::Bypass,
2591                bamboo_domain::PermissionMode::Default,
2592            ),
2593            &sink,
2594        );
2595
2596        let events = std::iter::from_fn(|| receiver.try_recv().ok()).collect::<Vec<_>>();
2597        assert!(events.iter().any(|event| {
2598            event["phase"] == "bootstrap"
2599                && event["sandbox"] == "danger-full-access"
2600                && event["approval_policy"] == "never"
2601                && event["policy_invocation"] == "danger-bypass"
2602                && event["requested_mode"] == "bypass"
2603                && event["effective_mode"] == "bypass"
2604                && event["executor_mapping"] == "codex_exec:approval_policy=never"
2605        }));
2606        assert!(events.iter().any(|event| {
2607            event["phase"] == "policy_warning"
2608                && event["level"] == "warning"
2609                && event["message"]
2610                    .as_str()
2611                    .is_some_and(|message| message.contains("DANGER"))
2612        }));
2613    }
2614
2615    #[test]
2616    fn auth_mode_matrix_is_explicit_isolated_and_secret_safe() {
2617        let credential = ScopedCredential {
2618            provider: "custom-provider".to_string(),
2619            api_key: "custom-secret-570".to_string(),
2620            base_url: None,
2621            provider_type: Some("openai".to_string()),
2622            credential_ref: Some("provider.custom.api_key".to_string()),
2623        };
2624
2625        let inherit = resolved_auth(Some("inherit"), None, None, &[], &[]).unwrap();
2626        assert_eq!(inherit.mode(), CodexAuthMode::Inherit);
2627        let mut inherit_executor = fixture_executor();
2628        inherit_executor.auth = inherit;
2629        let policy = default_policy(&inherit_executor);
2630        let inherit_command = inherit_executor.build_command(None, &policy, None).unwrap();
2631        assert!(command_env(&inherit_command, "CODEX_HOME").is_none());
2632        let inherit_args = inherit_command
2633            .as_std()
2634            .get_args()
2635            .map(|arg| arg.to_string_lossy().into_owned())
2636            .collect::<Vec<_>>();
2637        assert!(!inherit_args.iter().any(|arg| arg == "--ignore-rules"));
2638
2639        let api_key = resolved_auth(Some("api_key"), None, None, &[], &["OPENAI_API_KEY"]).unwrap();
2640        assert_eq!(api_key.mode(), CodexAuthMode::ApiKey);
2641        assert!(
2642            auth_error(resolved_auth(Some("api_key"), None, None, &[], &[]))
2643                .contains("explicit OPENAI_API_KEY")
2644        );
2645
2646        let custom = resolved_auth(
2647            Some("custom"),
2648            Some("https://provider.example/v1"),
2649            Some("provider.custom.api_key"),
2650            std::slice::from_ref(&credential),
2651            &[],
2652        )
2653        .unwrap();
2654        assert_eq!(custom.mode(), CodexAuthMode::Custom);
2655        let generated = custom.generated_config_toml().unwrap();
2656        assert!(generated.contains("model_provider = \"custom\""));
2657        assert!(generated.contains("base_url = \"https://provider.example/v1\""));
2658        assert!(generated.contains("env_key = \"BAMBOO_CODEX_PROVIDER_KEY\""));
2659        assert!(generated.contains("wire_api = \"responses\""));
2660        assert!(!generated.contains("custom-secret-570"));
2661
2662        let state = tempfile::tempdir().unwrap();
2663        let mut custom_executor = fixture_executor();
2664        custom_executor.state_dir = Some(state.path().to_path_buf());
2665        custom_executor.auth = custom;
2666        let policy = default_policy(&custom_executor);
2667        let custom_command = custom_executor.build_command(None, &policy, None).unwrap();
2668        assert_eq!(
2669            command_env(&custom_command, CODEX_PROVIDER_ENV).as_deref(),
2670            Some("custom-secret-570")
2671        );
2672        assert_eq!(
2673            command_env(&custom_command, "CODEX_HOME"),
2674            Some(
2675                state
2676                    .path()
2677                    .join("codex-home")
2678                    .to_string_lossy()
2679                    .into_owned()
2680            )
2681        );
2682        assert!(custom_command
2683            .as_std()
2684            .get_args()
2685            .any(|arg| arg == "--ignore-rules"));
2686
2687        let bamboo = resolved_auth(
2688            None,
2689            Some("http://127.0.0.1:9562/openai/v1"),
2690            None,
2691            &[],
2692            &[],
2693        )
2694        .unwrap();
2695        assert_eq!(bamboo.mode(), CodexAuthMode::Bamboo);
2696        let generated = bamboo.generated_config_toml().unwrap();
2697        assert!(generated.contains("model_provider = \"bamboo\""));
2698        assert!(generated.contains("http://127.0.0.1:9562/openai/v1"));
2699        assert!(!generated.contains("bcx1_"));
2700        let app_server_generated = bamboo
2701            .generated_app_server_config_toml(
2702                Path::new("/opt/bamboo/bin/bamboo"),
2703                Path::new("/private/state/codex-provider-token"),
2704            )
2705            .unwrap();
2706        assert!(app_server_generated.contains("[model_providers.bamboo.auth]"));
2707        assert!(app_server_generated.contains("command = \"/opt/bamboo/bin/bamboo\""));
2708        assert!(app_server_generated.contains("\"codex-provider-token\""));
2709        assert!(app_server_generated.contains("refresh_interval_ms = 1"));
2710        assert!(!app_server_generated.contains("env_key"));
2711        assert!(!app_server_generated.contains("bcx1_"));
2712        let mut bamboo_executor = fixture_executor();
2713        bamboo_executor.state_dir = Some(state.path().to_path_buf());
2714        bamboo_executor.auth = bamboo;
2715        let policy = default_policy(&bamboo_executor);
2716        assert!(bamboo_executor
2717            .build_command(None, &policy, None)
2718            .unwrap_err()
2719            .contains("per-run provider token"));
2720        let bamboo_command = bamboo_executor
2721            .build_command(Some("bcx1_per_run_570"), &policy, None)
2722            .unwrap();
2723        assert_eq!(
2724            command_env(&bamboo_command, CODEX_PROVIDER_ENV).as_deref(),
2725            Some("bcx1_per_run_570")
2726        );
2727    }
2728
2729    #[test]
2730    fn auth_matrix_rejects_ambiguous_or_unsafe_configuration() {
2731        for mode in ["inherit", "custom", "bamboo"] {
2732            let error = auth_error(resolved_auth(
2733                Some(mode),
2734                match mode {
2735                    "custom" => Some("https://provider.example/v1"),
2736                    "bamboo" => Some("http://127.0.0.1:9562/openai/v1"),
2737                    _ => None,
2738                },
2739                (mode == "custom").then_some("provider.missing.api_key"),
2740                &[],
2741                &["OPENAI_API_KEY"],
2742            ));
2743            assert!(
2744                error.contains("only be forwarded") || error.contains("did not resolve"),
2745                "mode {mode}: {error}"
2746            );
2747        }
2748        assert!(
2749            auth_error(resolved_auth(Some("future"), None, None, &[], &[]))
2750                .contains("unknown Codex auth mode")
2751        );
2752        assert!(auth_error(resolve_codex_auth_config(
2753            Some("inherit"),
2754            false,
2755            None,
2756            Some("chat".to_string()),
2757            None,
2758            &[],
2759            &[],
2760        ))
2761        .contains("requires responses"));
2762        assert!(auth_error(resolved_auth(
2763            Some("inherit"),
2764            None,
2765            None,
2766            &[],
2767            &["CODEX_HOME"],
2768        ))
2769        .contains("reserved variable"));
2770        assert!(auth_error(resolved_auth(
2771            Some("custom"),
2772            Some("https://user:secret@provider.example/v1?x=1"),
2773            Some("provider.custom.api_key"),
2774            &[],
2775            &[],
2776        ))
2777        .contains("must not contain credentials"));
2778    }
2779
2780    #[cfg(unix)]
2781    #[tokio::test]
2782    async fn isolated_home_removes_stale_login_and_writes_locked_secret_free_config() {
2783        use std::os::unix::fs::PermissionsExt as _;
2784
2785        let state = tempfile::tempdir().unwrap();
2786        let home = state.path().join("codex-home");
2787        std::fs::create_dir_all(&home).unwrap();
2788        std::fs::write(home.join("auth.json"), r#"{"tokens":"stale-secret"}"#).unwrap();
2789        let mut executor = fixture_executor();
2790        executor.state_dir = Some(state.path().to_path_buf());
2791        executor.auth = CodexAuthConfig {
2792            mode: CodexAuthMode::Custom,
2793            base_url: Some("https://provider.example/v1".to_string()),
2794            wire_api: "responses".to_string(),
2795            provider_key: Some("custom-secret-570".to_string()),
2796        };
2797
2798        executor.prepare_auth_home().await.unwrap();
2799
2800        assert!(!home.join("auth.json").exists());
2801        let config_path = home.join("config.toml");
2802        let config = std::fs::read_to_string(&config_path).unwrap();
2803        assert!(config.contains("https://provider.example/v1"));
2804        assert!(!config.contains("custom-secret-570"));
2805        assert_eq!(
2806            std::fs::metadata(&home).unwrap().permissions().mode() & 0o777,
2807            0o700
2808        );
2809        assert_eq!(
2810            std::fs::metadata(&config_path)
2811                .unwrap()
2812                .permissions()
2813                .mode()
2814                & 0o777,
2815            0o600
2816        );
2817
2818        // Simulate a prior activation mutating its reusable worker home. The
2819        // next run boundary must restore the managed config and remove auth.
2820        std::fs::write(home.join("auth.json"), r#"{"tokens":"late-secret"}"#).unwrap();
2821        std::fs::write(&config_path, "model_provider = \"attacker\"\n").unwrap();
2822        executor.prepare_auth_home().await.unwrap();
2823        assert!(!home.join("auth.json").exists());
2824        let restored = std::fs::read_to_string(&config_path).unwrap();
2825        assert!(restored.contains("model_provider = \"custom\""));
2826        assert!(!restored.contains("attacker"));
2827    }
2828
2829    #[tokio::test]
2830    async fn bounded_reader_accepts_ten_megabytes_and_rejects_more() {
2831        let allowed = vec![b'x'; MAX_STDOUT_LINE_BYTES];
2832        let mut bytes = allowed.clone();
2833        bytes.push(b'\n');
2834        let mut reader = BufReader::new(bytes.as_slice());
2835        assert_eq!(
2836            read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES)
2837                .await
2838                .unwrap()
2839                .unwrap()
2840                .len(),
2841            MAX_STDOUT_LINE_BYTES
2842        );
2843
2844        let mut too_large = vec![b'x'; MAX_STDOUT_LINE_BYTES + 1];
2845        too_large.push(b'\n');
2846        let mut reader = BufReader::new(too_large.as_slice());
2847        assert_eq!(
2848            read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES)
2849                .await
2850                .unwrap_err()
2851                .kind(),
2852            std::io::ErrorKind::InvalidData
2853        );
2854    }
2855
2856    #[test]
2857    fn usage_matches_real_codex_schema() {
2858        let usage = parse_usage(Some(&json!({
2859            "input_tokens": 13_460,
2860            "cached_input_tokens": 9_984,
2861            "output_tokens": 6,
2862            "reasoning_output_tokens": 0
2863        })));
2864        assert_eq!(usage.prompt_tokens, 13_460);
2865        assert_eq!(usage.completion_tokens, 6);
2866        assert_eq!(usage.total_tokens, 13_466);
2867    }
2868
2869    #[test]
2870    fn recorded_jsonl_fixtures_map_agent_and_command_events() {
2871        let cases = [
2872            (
2873                include_str!("../tests/fixtures/codex-cli/0.144.5-simple.jsonl"),
2874                "PONG",
2875                false,
2876                13_466,
2877            ),
2878            (
2879                include_str!("../tests/fixtures/codex-cli/0.144.5-command.jsonl"),
2880                "The current working directory is `/private/tmp/zenith-bamboo-569-codex-cli`.",
2881                true,
2882                27_108,
2883            ),
2884        ];
2885
2886        for (fixture, expected_final, expects_tool, total_tokens) in cases {
2887            let (state, events) = map_fixture(fixture);
2888            assert!(state.completed);
2889            assert_eq!(state.last_agent_message, expected_final);
2890            assert_eq!(state.usage.total_tokens, total_tokens);
2891            assert!(events.iter().any(|event| event["type"] == "token"));
2892            assert!(events.iter().any(|event| event["type"] == "complete"));
2893            assert_eq!(
2894                events.iter().any(|event| event["type"] == "tool_start"),
2895                expects_tool
2896            );
2897            if expects_tool {
2898                assert!(events.iter().any(|event| {
2899                    event["type"] == "tool_start" && event["tool_name"] == "Bash"
2900                }));
2901                assert!(events.iter().any(|event| event["type"] == "tool_token"));
2902                assert!(events.iter().any(|event| event["type"] == "tool_complete"));
2903            }
2904        }
2905    }
2906
2907    #[test]
2908    fn item_mapping_table_covers_non_command_codex_items() {
2909        let cases = [
2910            (
2911                json!({"id":"r1","type":"reasoning","text":"thinking"}),
2912                "reasoning_token",
2913                None,
2914            ),
2915            (
2916                json!({"id":"f1","type":"file_change","changes":[{"path":"a.txt","kind":"add"}]}),
2917                "tool_start",
2918                Some("ApplyPatch"),
2919            ),
2920            (
2921                json!({"id":"m1","type":"mcp_tool_call","server":"files","tool":"read","arguments":{"path":"a.txt"},"result":"ok"}),
2922                "tool_start",
2923                Some("files::read"),
2924            ),
2925            (
2926                json!({"id":"w1","type":"web_search","query":"Bamboo","result":[]}),
2927                "tool_start",
2928                Some("WebSearch"),
2929            ),
2930            (
2931                json!({"id":"t1","type":"todo_list","items":[]}),
2932                "runner_progress",
2933                None,
2934            ),
2935        ];
2936
2937        for (item, expected_type, expected_tool) in cases {
2938            let (sink, mut rx) = EventSink::channel();
2939            let mut state = RunState::default();
2940            handle_item("completed", &item, &sink, &mut state);
2941            let first = rx.try_recv().expect("mapping emitted an event");
2942            assert_eq!(first["type"], expected_type, "item: {item}");
2943            if let Some(tool) = expected_tool {
2944                assert_eq!(first["tool_name"], tool, "item: {item}");
2945                assert!(rx
2946                    .try_recv()
2947                    .is_ok_and(|event| event["type"] == "tool_complete"));
2948            }
2949        }
2950    }
2951
2952    #[test]
2953    fn unknown_event_and_item_types_are_tolerated() {
2954        let executor = fixture_executor();
2955        let (sink, mut rx) = EventSink::channel();
2956        let mut state = RunState::default();
2957        let policy = default_policy(&executor);
2958        executor.handle_event(
2959            json!({"type":"future.event","payload":{"schema":2}}),
2960            &policy,
2961            &sink,
2962            &mut state,
2963        );
2964        executor.handle_event(
2965            json!({"type":"item.completed","item":{"id":"x","type":"future_item"}}),
2966            &policy,
2967            &sink,
2968            &mut state,
2969        );
2970        assert!(rx.try_recv().is_err());
2971        assert!(!state.completed);
2972        assert!(state.failure.is_none());
2973    }
2974
2975    #[test]
2976    fn sandbox_denial_report_without_cli_tool_item_becomes_tool_error() {
2977        let executor = fixture_executor();
2978        let policy = executor.permissions.effective(false, false);
2979        let (sink, mut receiver) = EventSink::channel();
2980        let mut state = RunState::default();
2981        executor.handle_event(
2982            json!({
2983                "type": "item.completed",
2984                "item": {
2985                    "id": "denied-report",
2986                    "type": "agent_message",
2987                    "text": "Command exited with status 1. Sandbox failure: outside write: Operation not permitted"
2988                }
2989            }),
2990            &policy,
2991            &sink,
2992            &mut state,
2993        );
2994
2995        let events = std::iter::from_fn(|| receiver.try_recv().ok()).collect::<Vec<_>>();
2996        assert!(events.iter().any(|event| event["type"] == "token"));
2997        assert!(events.iter().any(|event| {
2998            event["type"] == "tool_error"
2999                && event["tool_call_id"] == "denied-report-sandbox-denial"
3000                && event["error"]
3001                    .as_str()
3002                    .is_some_and(|error| error.contains("Operation not permitted"))
3003        }));
3004    }
3005
3006    #[cfg(unix)]
3007    #[test]
3008    fn provider_token_reader_checks_open_descriptor_and_rejects_symlinks() {
3009        use std::os::unix::fs::{symlink, PermissionsExt as _};
3010
3011        let root = tempfile::tempdir().unwrap();
3012        let token = root.path().join("token");
3013        std::fs::write(&token, "  bcx1_short_lived  \n").unwrap();
3014        std::fs::set_permissions(&token, std::fs::Permissions::from_mode(0o600)).unwrap();
3015        assert_eq!(
3016            read_codex_provider_token(&token).unwrap(),
3017            "bcx1_short_lived"
3018        );
3019
3020        let link = root.path().join("token-link");
3021        symlink(&token, &link).unwrap();
3022        assert!(read_codex_provider_token(&link).is_err());
3023
3024        std::fs::set_permissions(&token, std::fs::Permissions::from_mode(0o640)).unwrap();
3025        assert!(read_codex_provider_token(&token)
3026            .unwrap_err()
3027            .contains("group/other"));
3028    }
3029
3030    #[cfg(unix)]
3031    mod unix_process_tests {
3032        use super::*;
3033        use std::io::Write as _;
3034        use std::os::unix::fs::PermissionsExt;
3035
3036        fn write_stub(dir: &Path, body: &str) -> PathBuf {
3037            let path = dir.join("codex");
3038            let mut file = std::fs::File::create(&path).unwrap();
3039            writeln!(file, "#!/bin/sh").unwrap();
3040            file.write_all(body.as_bytes()).unwrap();
3041            let mut permissions = std::fs::metadata(&path).unwrap().permissions();
3042            permissions.set_mode(0o755);
3043            std::fs::set_permissions(&path, permissions).unwrap();
3044            path
3045        }
3046
3047        fn executor(binary: PathBuf, workspace: &Path) -> CodexExecutor {
3048            CodexExecutor {
3049                binary,
3050                version: "codex-cli 0.144.5".to_string(),
3051                model: None,
3052                permissions: permissions(None, None, false, false, None, false, true),
3053                workspace: Some(workspace.to_string_lossy().into_owned()),
3054                state_dir: None,
3055                forward_env: Vec::new(),
3056                auth: CodexAuthConfig::inherit(),
3057            }
3058        }
3059
3060        fn run_spec(assignment: &str) -> RunSpec {
3061            RunSpec {
3062                assignment: assignment.to_string(),
3063                logical_session: None,
3064                project_id: None,
3065                reasoning_effort: None,
3066                permission_policy: None,
3067                messages: Vec::new(),
3068                activation_run_id: None,
3069                initial_session_messages: Vec::new(),
3070                secrets: Default::default(),
3071            }
3072        }
3073
3074        fn run_spec_with_messages(assignment: &str, messages: Vec<Value>) -> RunSpec {
3075            RunSpec {
3076                assignment: assignment.to_string(),
3077                logical_session: None,
3078                project_id: None,
3079                reasoning_effort: None,
3080                permission_policy: None,
3081                messages,
3082                activation_run_id: None,
3083                initial_session_messages: Vec::new(),
3084                secrets: Default::default(),
3085            }
3086        }
3087
3088        fn message(role: &str, content: &str) -> Value {
3089            json!({"role": role, "content": content})
3090        }
3091
3092        #[tokio::test]
3093        async fn preflight_accepts_current_surface_and_rejects_old_or_missing_binary() {
3094            let dir = tempfile::tempdir().unwrap();
3095            let current = write_stub(
3096                dir.path(),
3097                r#"
3098case "$*" in
3099  "--version") echo 'codex-cli 0.144.5' ;;
3100  "exec --help") echo '--json --output-last-message --config --sandbox --dangerously-bypass-approvals-and-sandbox prompt from stdin' ;;
3101  "exec resume --help") echo '--json' ;;
3102  *) exit 2 ;;
3103esac
3104"#,
3105            );
3106            let checked = CodexExecutor::new(
3107                Some(current.to_string_lossy().into_owned()),
3108                None,
3109                None,
3110                None,
3111                Vec::new(),
3112                CodexAuthConfig::inherit(),
3113                default_permissions(),
3114            )
3115            .await
3116            .unwrap();
3117            assert_eq!(checked.version, "codex-cli 0.144.5");
3118
3119            let old_dir = tempfile::tempdir().unwrap();
3120            let old = write_stub(
3121                old_dir.path(),
3122                r#"
3123if [ "$1" = "--version" ]; then echo 'codex-cli 0.143.9'; else exit 2; fi
3124"#,
3125            );
3126            let old_error = CodexExecutor::new(
3127                Some(old.to_string_lossy().into_owned()),
3128                None,
3129                None,
3130                None,
3131                Vec::new(),
3132                CodexAuthConfig::inherit(),
3133                default_permissions(),
3134            )
3135            .await
3136            .err()
3137            .expect("old version rejected");
3138            assert!(old_error.contains("too old"));
3139            assert!(old_error.contains(">= 0.144.0"));
3140
3141            let missing_error = CodexExecutor::new(
3142                Some("/definitely/missing/codex".to_string()),
3143                None,
3144                None,
3145                None,
3146                Vec::new(),
3147                CodexAuthConfig::inherit(),
3148                default_permissions(),
3149            )
3150            .await
3151            .err()
3152            .expect("missing binary rejected");
3153            assert!(missing_error.contains("npm i -g @openai/codex"));
3154            assert!(missing_error.contains("brew install codex"));
3155            assert!(missing_error.contains("codex_binary"));
3156        }
3157
3158        #[tokio::test]
3159        async fn clean_completion_and_nonzero_exit_have_distinct_outcomes() {
3160            let workspace = tempfile::tempdir().unwrap();
3161            let ok_dir = tempfile::tempdir().unwrap();
3162            let ok = write_stub(
3163                ok_dir.path(),
3164                r#"
3165read -r prompt
3166echo '{"type":"thread.started","thread_id":"stub-thread"}'
3167echo '{"type":"turn.started"}'
3168echo '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"PONG"}}'
3169echo '{"type":"turn.completed","usage":{"input_tokens":3,"output_tokens":1}}'
3170"#,
3171            );
3172            let (sink, _rx) = EventSink::channel();
3173            let outcome = executor(ok, workspace.path())
3174                .run(
3175                    run_spec("reply PONG"),
3176                    sink,
3177                    SteerInbox::disconnected(),
3178                    CancellationToken::new(),
3179                )
3180                .await;
3181            assert_eq!(outcome.status, TerminalStatus::Completed);
3182            assert_eq!(outcome.result.as_deref(), Some("PONG"));
3183
3184            let fail_dir = tempfile::tempdir().unwrap();
3185            let fail = write_stub(
3186                fail_dir.path(),
3187                r#"
3188read -r prompt
3189echo 'credential lookup failed' >&2
3190exit 7
3191"#,
3192            );
3193            let (sink, _rx) = EventSink::channel();
3194            let outcome = executor(fail, workspace.path())
3195                .run(
3196                    run_spec("fail"),
3197                    sink,
3198                    SteerInbox::disconnected(),
3199                    CancellationToken::new(),
3200                )
3201                .await;
3202            assert_eq!(outcome.status, TerminalStatus::Error);
3203            let error = outcome.error.unwrap();
3204            assert!(error.contains("status exit status: 7"), "{error}");
3205            assert!(error.contains("credential lookup failed"), "{error}");
3206        }
3207
3208        #[tokio::test]
3209        async fn explicit_deny_fails_before_codex_exec_spawn_without_leaking_rule_resource() {
3210            let workspace = tempfile::tempdir().unwrap();
3211            let bin_dir = tempfile::tempdir().unwrap();
3212            let marker = bin_dir.path().join("spawned");
3213            let bin = write_stub(
3214                bin_dir.path(),
3215                r#"
3216DIR=$(cd "$(dirname "$0")" && pwd)
3217: > "$DIR/spawned"
3218exit 0
3219"#,
3220            );
3221            let secret_resource = "TOP_SECRET_CODEX_EXEC_DENY_RESOURCE";
3222            let mut policy = bamboo_tools::permission::SerializablePermissionConfig::default();
3223            policy
3224                .whitelist
3225                .push(bamboo_tools::permission::PermissionRule::new(
3226                    bamboo_tools::permission::PermissionType::ExecuteCommand,
3227                    secret_resource,
3228                    false,
3229                ));
3230            let mut spec = run_spec("must fail closed");
3231            spec.permission_policy = Some(bamboo_subagent::proto::PermissionPolicyContext {
3232                revision: 23,
3233                requested_mode: "auto".to_string(),
3234                effective_mode: "auto".to_string(),
3235                bypass_permissions: false,
3236                auto_approve_permissions: true,
3237                session_id: "codex-exec-explicit-deny".to_string(),
3238                workspace_path: Some(workspace.path().to_string_lossy().into_owned()),
3239                inherit_session_grants: false,
3240                policy: serde_json::to_value(policy).unwrap(),
3241            });
3242            let (sink, mut rx) = EventSink::channel();
3243
3244            let outcome = executor(bin, workspace.path())
3245                .run(
3246                    spec,
3247                    sink,
3248                    SteerInbox::disconnected(),
3249                    CancellationToken::new(),
3250                )
3251                .await;
3252
3253            assert_eq!(outcome.status, TerminalStatus::Error);
3254            assert!(outcome
3255                .error
3256                .as_deref()
3257                .is_some_and(|error| error.contains("explicit-deny")));
3258            assert!(!marker.exists(), "Codex exec process must not be spawned");
3259            let events = std::iter::from_fn(|| rx.try_recv().ok()).collect::<Vec<_>>();
3260            assert!(events.iter().any(|event| {
3261                event["type"] == "permission_posture_activated"
3262                    && event["executor_mapping"] == "codex_exec:blocked_explicit_deny"
3263            }));
3264            assert!(events.iter().any(|event| event["type"] == "error"));
3265            assert!(
3266                !serde_json::to_string(&events)
3267                    .unwrap()
3268                    .contains(secret_resource),
3269                "deny resources must not be emitted in audit/error events"
3270            );
3271        }
3272
3273        #[tokio::test]
3274        async fn spawn_uses_stdin_safe_defaults_and_last_message_fallback() {
3275            let workspace = tempfile::tempdir().unwrap();
3276            let bin_dir = tempfile::tempdir().unwrap();
3277            let state_dir = tempfile::tempdir().unwrap();
3278            let bin = write_stub(
3279                bin_dir.path(),
3280                r#"
3281DIR=$(cd "$(dirname "$0")" && pwd)
3282: > "$DIR/argv.txt"
3283OUT=''
3284while [ "$#" -gt 0 ]; do
3285  printf '%s\n' "$1" >> "$DIR/argv.txt"
3286  if [ "$1" = '--output-last-message' ]; then
3287    shift
3288    OUT="$1"
3289    printf '%s\n' "$1" >> "$DIR/argv.txt"
3290  fi
3291  shift
3292done
3293IFS= read -r prompt
3294printf '%s\n' "$prompt" > "$DIR/stdin.txt"
3295printf 'PONG\n' > "$OUT"
3296echo '{"type":"thread.started","thread_id":"fallback-thread"}'
3297echo '{"type":"turn.started"}'
3298echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
3299"#,
3300            );
3301            let mut exec = executor(bin, workspace.path());
3302            exec.state_dir = Some(state_dir.path().to_path_buf());
3303            let assignment = "a private prompt that must not appear in argv";
3304            let (sink, _rx) = EventSink::channel();
3305            let outcome = exec
3306                .run(
3307                    run_spec(assignment),
3308                    sink,
3309                    SteerInbox::disconnected(),
3310                    CancellationToken::new(),
3311                )
3312                .await;
3313            assert_eq!(outcome.status, TerminalStatus::Completed);
3314            assert_eq!(outcome.result.as_deref(), Some("PONG"));
3315            assert_eq!(
3316                std::fs::read_to_string(bin_dir.path().join("stdin.txt"))
3317                    .unwrap()
3318                    .trim(),
3319                assignment
3320            );
3321            let argv = std::fs::read_to_string(bin_dir.path().join("argv.txt")).unwrap();
3322            assert!(
3323                !argv.contains(assignment),
3324                "prompt leaked into argv: {argv}"
3325            );
3326            for required in [
3327                "exec",
3328                "--json",
3329                "--cd",
3330                "--skip-git-repo-check",
3331                "--sandbox",
3332                "workspace-write",
3333                "--config",
3334                "approval_policy=\"never\"",
3335                "--output-last-message",
3336                "-",
3337            ] {
3338                assert!(
3339                    argv.lines().any(|arg| arg == required),
3340                    "missing {required}: {argv}"
3341                );
3342            }
3343            assert!(!argv.lines().any(|arg| arg == "--ignore-rules"));
3344            assert!(!argv.lines().any(|arg| arg == "--ignore-user-config"));
3345        }
3346
3347        #[tokio::test]
3348        async fn thread_state_is_captured_resumed_recaptured_and_cleared_for_fresh_run() {
3349            let workspace = tempfile::tempdir().unwrap();
3350            let bin_dir = tempfile::tempdir().unwrap();
3351            let state_dir = tempfile::tempdir().unwrap();
3352            let bin = write_stub(
3353                bin_dir.path(),
3354                r#"
3355DIR=$(cd "$(dirname "$0")" && pwd)
3356N=$(cat "$DIR/count" 2>/dev/null || echo 0)
3357N=$((N+1))
3358echo "$N" > "$DIR/count"
3359printf '%s\n' "$@" > "$DIR/argv-$N.txt"
3360IFS= read -r prompt
3361printf '%s\n' "$prompt" > "$DIR/stdin-$N.txt"
3362case "$N" in
3363  1) THREAD='thread-one'; ANSWER='first' ;;
3364  2) THREAD='thread-two'; ANSWER='second' ;;
3365  *) THREAD=''; ANSWER='third' ;;
3366esac
3367if [ -n "$THREAD" ]; then
3368  printf '{"type":"thread.started","thread_id":"%s"}\n' "$THREAD"
3369fi
3370echo '{"type":"turn.started"}'
3371printf '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"%s"}}\n' "$ANSWER"
3372echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
3373"#,
3374            );
3375            let mut exec = executor(bin, workspace.path());
3376            exec.state_dir = Some(state_dir.path().to_path_buf());
3377            let state_path = state_dir.path().join(CODEX_SESSION_STATE_FILE);
3378
3379            let (sink, _rx) = EventSink::channel();
3380            let first = exec
3381                .run(
3382                    run_spec("remember amber-572"),
3383                    sink,
3384                    SteerInbox::disconnected(),
3385                    CancellationToken::new(),
3386                )
3387                .await;
3388            assert_eq!(first.status, TerminalStatus::Completed);
3389            assert!(!std::fs::read_to_string(bin_dir.path().join("argv-1.txt"))
3390                .unwrap()
3391                .lines()
3392                .any(|argument| argument == "resume"));
3393            let state: CodexSessionState =
3394                serde_json::from_slice(&std::fs::read(&state_path).unwrap()).unwrap();
3395            assert_eq!(state.thread_id, "thread-one");
3396            assert_eq!(state.workspace, exec.workspace);
3397            assert_eq!(state.codex_home_mode, "inherit");
3398
3399            let (sink, _rx) = EventSink::channel();
3400            let second = exec
3401                .run(
3402                    run_spec_with_messages(
3403                        "what nonce?",
3404                        vec![
3405                            message("user", "remember amber-572"),
3406                            message("assistant", "stored"),
3407                            message("user", "what nonce?"),
3408                        ],
3409                    ),
3410                    sink,
3411                    SteerInbox::disconnected(),
3412                    CancellationToken::new(),
3413                )
3414                .await;
3415            assert_eq!(second.status, TerminalStatus::Completed);
3416            let argv = std::fs::read_to_string(bin_dir.path().join("argv-2.txt")).unwrap();
3417            let arguments = argv.lines().collect::<Vec<_>>();
3418            let resume_index = arguments
3419                .iter()
3420                .position(|argument| *argument == "resume")
3421                .expect("resume subcommand");
3422            assert_eq!(arguments.get(resume_index + 1), Some(&"thread-one"));
3423            assert_eq!(arguments.last(), Some(&"-"));
3424            assert_eq!(
3425                std::fs::read_to_string(bin_dir.path().join("stdin-2.txt"))
3426                    .unwrap()
3427                    .trim(),
3428                "what nonce?"
3429            );
3430            let state: CodexSessionState =
3431                serde_json::from_slice(&std::fs::read(&state_path).unwrap()).unwrap();
3432            assert_eq!(state.thread_id, "thread-two");
3433
3434            let (sink, _rx) = EventSink::channel();
3435            let third = exec
3436                .run(
3437                    run_spec("fresh rerun"),
3438                    sink,
3439                    SteerInbox::disconnected(),
3440                    CancellationToken::new(),
3441                )
3442                .await;
3443            assert_eq!(third.status, TerminalStatus::Completed);
3444            assert!(!state_path.exists());
3445            assert!(!std::fs::read_to_string(bin_dir.path().join("argv-3.txt"))
3446                .unwrap()
3447                .lines()
3448                .any(|argument| argument == "resume"));
3449        }
3450
3451        #[tokio::test]
3452        async fn missing_or_mismatched_state_uses_bounded_history_rehydration() {
3453            let workspace = tempfile::tempdir().unwrap();
3454            let bin_dir = tempfile::tempdir().unwrap();
3455            let state_dir = tempfile::tempdir().unwrap();
3456            let bin = write_stub(
3457                bin_dir.path(),
3458                r#"
3459DIR=$(cd "$(dirname "$0")" && pwd)
3460printf '%s\n' "$@" > "$DIR/argv.txt"
3461cat > "$DIR/stdin.txt"
3462echo '{"type":"thread.started","thread_id":"fallback-thread"}'
3463echo '{"type":"turn.started"}'
3464echo '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"ok"}}'
3465echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
3466"#,
3467            );
3468            let mut exec = executor(bin, workspace.path());
3469            exec.state_dir = Some(state_dir.path().to_path_buf());
3470            write_json_atomic(
3471                &state_dir.path().join(CODEX_SESSION_STATE_FILE),
3472                &CodexSessionState {
3473                    thread_id: "unusable-thread".to_string(),
3474                    workspace: exec.workspace.clone(),
3475                    codex_home_mode: "isolated".to_string(),
3476                    updated_at: Utc::now(),
3477                },
3478            )
3479            .await
3480            .unwrap();
3481
3482            let (sink, _rx) = EventSink::channel();
3483            let outcome = exec
3484                .run(
3485                    run_spec_with_messages(
3486                        "continue",
3487                        vec![
3488                            message("user", "old question"),
3489                            message("assistant", "old answer"),
3490                            message("user", "continue"),
3491                        ],
3492                    ),
3493                    sink,
3494                    SteerInbox::disconnected(),
3495                    CancellationToken::new(),
3496                )
3497                .await;
3498            assert_eq!(outcome.status, TerminalStatus::Completed);
3499            let argv = std::fs::read_to_string(bin_dir.path().join("argv.txt")).unwrap();
3500            assert!(!argv.lines().any(|argument| argument == "resume"));
3501            assert!(!argv.contains("unusable-thread"));
3502            let body = std::fs::read_to_string(bin_dir.path().join("stdin.txt")).unwrap();
3503            assert!(body.contains("## Prior conversation (rehydrated)"));
3504            assert!(body.contains("old question"));
3505            assert!(body.contains("old answer"));
3506            assert!(body.contains("## Current task"));
3507            assert_eq!(body.matches("continue").count(), 1);
3508
3509            let state: CodexSessionState = serde_json::from_slice(
3510                &std::fs::read(state_dir.path().join(CODEX_SESSION_STATE_FILE)).unwrap(),
3511            )
3512            .unwrap();
3513            assert_eq!(state.thread_id, "fallback-thread");
3514
3515            let workspace_mismatch = CodexSessionState {
3516                thread_id: "wrong-workspace".to_string(),
3517                workspace: Some("/different/workspace".to_string()),
3518                codex_home_mode: "inherit".to_string(),
3519                updated_at: Utc::now(),
3520            };
3521            write_json_atomic(
3522                &state_dir.path().join(CODEX_SESSION_STATE_FILE),
3523                &workspace_mismatch,
3524            )
3525            .await
3526            .unwrap();
3527            assert_eq!(exec.resolve_resume_id().await, None);
3528        }
3529
3530        #[tokio::test]
3531        async fn invalid_resume_retries_once_fresh_with_rehydrated_history() {
3532            let workspace = tempfile::tempdir().unwrap();
3533            let bin_dir = tempfile::tempdir().unwrap();
3534            let state_dir = tempfile::tempdir().unwrap();
3535            let bin = write_stub(
3536                bin_dir.path(),
3537                r#"
3538DIR=$(cd "$(dirname "$0")" && pwd)
3539N=$(cat "$DIR/count" 2>/dev/null || echo 0)
3540N=$((N+1))
3541echo "$N" > "$DIR/count"
3542printf '%s\n' "$@" > "$DIR/argv-$N.txt"
3543cat > "$DIR/stdin-$N.txt"
3544case " $* " in
3545  *' resume dead-thread '*)
3546    echo 'thread not found' >&2
3547    exit 1
3548    ;;
3549  *)
3550    echo '{"type":"thread.started","thread_id":"recovered-thread"}'
3551    echo '{"type":"turn.started"}'
3552    echo '{"type":"item.completed","item":{"id":"answer","type":"agent_message","text":"recovered"}}'
3553    echo '{"type":"turn.completed","usage":{"input_tokens":1,"output_tokens":1}}'
3554    ;;
3555esac
3556"#,
3557            );
3558            let mut exec = executor(bin, workspace.path());
3559            exec.state_dir = Some(state_dir.path().to_path_buf());
3560            write_json_atomic(
3561                &state_dir.path().join(CODEX_SESSION_STATE_FILE),
3562                &CodexSessionState {
3563                    thread_id: "dead-thread".to_string(),
3564                    workspace: exec.workspace.clone(),
3565                    codex_home_mode: "inherit".to_string(),
3566                    updated_at: Utc::now(),
3567                },
3568            )
3569            .await
3570            .unwrap();
3571
3572            let (sink, mut events) = EventSink::channel();
3573            let outcome = exec
3574                .run(
3575                    run_spec_with_messages(
3576                        "please continue",
3577                        vec![
3578                            message("user", "prior nonce amber-572"),
3579                            message("assistant", "stored"),
3580                            message("user", "please continue"),
3581                        ],
3582                    ),
3583                    sink,
3584                    SteerInbox::disconnected(),
3585                    CancellationToken::new(),
3586                )
3587                .await;
3588            assert_eq!(outcome.status, TerminalStatus::Completed);
3589            assert_eq!(outcome.result.as_deref(), Some("recovered"));
3590            assert_eq!(
3591                std::fs::read_to_string(bin_dir.path().join("count"))
3592                    .unwrap()
3593                    .trim(),
3594                "2"
3595            );
3596            let first_argv = std::fs::read_to_string(bin_dir.path().join("argv-1.txt")).unwrap();
3597            assert!(first_argv.lines().any(|argument| argument == "resume"));
3598            assert!(first_argv.lines().any(|argument| argument == "dead-thread"));
3599            let second_argv = std::fs::read_to_string(bin_dir.path().join("argv-2.txt")).unwrap();
3600            assert!(!second_argv.lines().any(|argument| argument == "resume"));
3601            let fallback = std::fs::read_to_string(bin_dir.path().join("stdin-2.txt")).unwrap();
3602            assert!(fallback.contains("## Prior conversation (rehydrated)"));
3603            assert!(fallback.contains("prior nonce amber-572"));
3604            assert!(std::iter::from_fn(|| events.try_recv().ok()).any(|event| {
3605                event["type"] == "runner_progress" && event["phase"] == "resume_fallback"
3606            }));
3607
3608            let state: CodexSessionState = serde_json::from_slice(
3609                &std::fs::read(state_dir.path().join(CODEX_SESSION_STATE_FILE)).unwrap(),
3610            )
3611            .unwrap();
3612            assert_eq!(state.thread_id, "recovered-thread");
3613            assert!(!std::fs::read_dir(state_dir.path())
3614                .unwrap()
3615                .filter_map(Result::ok)
3616                .any(|entry| entry.file_name().to_string_lossy().contains(".tmp.")));
3617        }
3618
3619        #[tokio::test]
3620        async fn resume_failure_after_turn_progress_is_not_retried() {
3621            let workspace = tempfile::tempdir().unwrap();
3622            let bin_dir = tempfile::tempdir().unwrap();
3623            let state_dir = tempfile::tempdir().unwrap();
3624            let bin = write_stub(
3625                bin_dir.path(),
3626                r#"
3627DIR=$(cd "$(dirname "$0")" && pwd)
3628N=$(cat "$DIR/count" 2>/dev/null || echo 0)
3629N=$((N+1))
3630echo "$N" > "$DIR/count"
3631cat >/dev/null
3632echo '{"type":"thread.started","thread_id":"progressed-thread"}'
3633echo '{"type":"turn.started"}'
3634echo '{"type":"turn.failed","error":{"message":"model failed"}}'
3635exit 1
3636"#,
3637            );
3638            let mut exec = executor(bin, workspace.path());
3639            exec.state_dir = Some(state_dir.path().to_path_buf());
3640            write_json_atomic(
3641                &state_dir.path().join(CODEX_SESSION_STATE_FILE),
3642                &CodexSessionState {
3643                    thread_id: "resumable-thread".to_string(),
3644                    workspace: exec.workspace.clone(),
3645                    codex_home_mode: "inherit".to_string(),
3646                    updated_at: Utc::now(),
3647                },
3648            )
3649            .await
3650            .unwrap();
3651
3652            let (sink, mut events) = EventSink::channel();
3653            let outcome = exec
3654                .run(
3655                    run_spec_with_messages(
3656                        "continue",
3657                        vec![message("user", "earlier"), message("user", "continue")],
3658                    ),
3659                    sink,
3660                    SteerInbox::disconnected(),
3661                    CancellationToken::new(),
3662                )
3663                .await;
3664
3665            assert_eq!(outcome.status, TerminalStatus::Error);
3666            assert_eq!(
3667                std::fs::read_to_string(bin_dir.path().join("count"))
3668                    .unwrap()
3669                    .trim(),
3670                "1"
3671            );
3672            assert!(!std::iter::from_fn(|| events.try_recv().ok()).any(|event| {
3673                event["type"] == "runner_progress" && event["phase"] == "resume_fallback"
3674            }));
3675        }
3676
3677        #[tokio::test]
3678        async fn failed_fallback_is_not_retried_a_second_time() {
3679            let workspace = tempfile::tempdir().unwrap();
3680            let bin_dir = tempfile::tempdir().unwrap();
3681            let state_dir = tempfile::tempdir().unwrap();
3682            let bin = write_stub(
3683                bin_dir.path(),
3684                r#"
3685DIR=$(cd "$(dirname "$0")" && pwd)
3686N=$(cat "$DIR/count" 2>/dev/null || echo 0)
3687N=$((N+1))
3688echo "$N" > "$DIR/count"
3689cat >/dev/null
3690echo "attempt $N failed" >&2
3691exit 1
3692"#,
3693            );
3694            let mut exec = executor(bin, workspace.path());
3695            exec.state_dir = Some(state_dir.path().to_path_buf());
3696            let state_path = state_dir.path().join(CODEX_SESSION_STATE_FILE);
3697            write_json_atomic(
3698                &state_path,
3699                &CodexSessionState {
3700                    thread_id: "dead-thread".to_string(),
3701                    workspace: exec.workspace.clone(),
3702                    codex_home_mode: "inherit".to_string(),
3703                    updated_at: Utc::now(),
3704                },
3705            )
3706            .await
3707            .unwrap();
3708
3709            let (sink, mut events) = EventSink::channel();
3710            let outcome = exec
3711                .run(
3712                    run_spec_with_messages(
3713                        "continue",
3714                        vec![message("user", "earlier"), message("user", "continue")],
3715                    ),
3716                    sink,
3717                    SteerInbox::disconnected(),
3718                    CancellationToken::new(),
3719                )
3720                .await;
3721
3722            assert_eq!(outcome.status, TerminalStatus::Error);
3723            assert_eq!(
3724                std::fs::read_to_string(bin_dir.path().join("count"))
3725                    .unwrap()
3726                    .trim(),
3727                "2"
3728            );
3729            assert_eq!(
3730                std::iter::from_fn(|| events.try_recv().ok())
3731                    .filter(|event| {
3732                        event["type"] == "runner_progress" && event["phase"] == "resume_fallback"
3733                    })
3734                    .count(),
3735                1
3736            );
3737            assert!(!state_path.exists());
3738        }
3739
3740        #[tokio::test]
3741        async fn cancellation_kills_the_entire_process_group() {
3742            let workspace = tempfile::tempdir().unwrap();
3743            let bin_dir = tempfile::tempdir().unwrap();
3744            let bin = write_stub(
3745                bin_dir.path(),
3746                r#"
3747DIR=$(cd "$(dirname "$0")" && pwd)
3748read -r prompt
3749echo '{"type":"thread.started","thread_id":"cancel-thread"}'
3750(trap '' TERM; sleep 30) &
3751echo $! > "$DIR/grandchild.pid"
3752wait
3753"#,
3754            );
3755            let pid_path = bin_dir.path().join("grandchild.pid");
3756            let (sink, _rx) = EventSink::channel();
3757            let cancel = CancellationToken::new();
3758            let cancel_for_run = cancel.clone();
3759            let run = tokio::spawn(async move {
3760                executor(bin, workspace.path())
3761                    .run(
3762                        run_spec("wait"),
3763                        sink,
3764                        SteerInbox::disconnected(),
3765                        cancel_for_run,
3766                    )
3767                    .await
3768            });
3769
3770            for _ in 0..500 {
3771                if pid_path.exists() {
3772                    break;
3773                }
3774                tokio::time::sleep(Duration::from_millis(10)).await;
3775            }
3776            assert!(pid_path.exists(), "grandchild pid was recorded");
3777            let grandchild_pid: libc::pid_t = std::fs::read_to_string(&pid_path)
3778                .unwrap()
3779                .trim()
3780                .parse()
3781                .unwrap();
3782            cancel.cancel();
3783            let outcome = tokio::time::timeout(Duration::from_secs(10), run)
3784                .await
3785                .expect("cancel completed within TERM/KILL bound")
3786                .unwrap();
3787            assert_eq!(outcome.status, TerminalStatus::Cancelled);
3788
3789            for _ in 0..100 {
3790                // SAFETY: signal 0 only probes existence and does not signal.
3791                if unsafe { libc::kill(grandchild_pid, 0) } == -1 {
3792                    break;
3793                }
3794                tokio::time::sleep(Duration::from_millis(10)).await;
3795            }
3796            // SAFETY: signal 0 only probes existence and does not signal.
3797            assert_eq!(unsafe { libc::kill(grandchild_pid, 0) }, -1);
3798        }
3799    }
3800}