1use 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
33pub 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#[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 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
210pub 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#[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
319pub 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 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 (
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#[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 unsafe { libc::geteuid() == 0 }
736}
737
738#[cfg(not(unix))]
739fn running_as_root() -> bool {
740 false
741}
742
743pub 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 if unsafe { libc::kill(grandchild_pid, 0) } == -1 {
3792 break;
3793 }
3794 tokio::time::sleep(Duration::from_millis(10)).await;
3795 }
3796 assert_eq!(unsafe { libc::kill(grandchild_pid, 0) }, -1);
3798 }
3799 }
3800}