Skip to main content

bamboo_agent/
codex_app_server_executor.rs

1//! Long-lived Codex app-server executor with Bamboo approval relay.
2//!
3//! The app-server process is retained by a warm actor worker, while logical
4//! Codex threads are keyed by Bamboo session id. Server-to-client approval
5//! requests are relayed through `HostBridge` and fail closed after 300 seconds.
6
7use std::collections::{HashMap, HashSet};
8use std::io::Write as _;
9use std::path::{Path, PathBuf};
10use std::process::Stdio;
11use std::sync::atomic::{AtomicU64, Ordering};
12use std::sync::Arc;
13use std::time::Duration;
14
15use async_trait::async_trait;
16use chrono::{DateTime, Utc};
17use serde::{Deserialize, Serialize};
18use serde_json::{json, Value};
19use tokio::io::{AsyncWriteExt, BufReader};
20use tokio::process::{Child, Command};
21use tokio::sync::{mpsc, oneshot, Mutex};
22use tokio_util::sync::CancellationToken;
23
24use bamboo_agent_core::{AgentEvent, TokenUsage, ToolResult};
25use bamboo_subagent::codex_discovery::discover_codex_app_server;
26use bamboo_subagent::executor::{ChildExecutor, ChildOutcome, EventSink, HostBridge, SteerInbox};
27use bamboo_subagent::executor_util::{build_rehydrated_turn, write_json_atomic};
28use bamboo_subagent::proto::RunSpec;
29
30use crate::codex_cli_executor::{
31    read_bounded_line, terminate_child, CodexAuthConfig, CodexAuthMode, CodexPermissionConfig,
32};
33
34const MAX_STDOUT_LINE_BYTES: usize = 10 * 1024 * 1024;
35const STDERR_TAIL_BYTES: usize = 16 * 1024;
36const TOOL_RESULT_TRUNCATE_CHARS: usize = 20_000;
37const REQUEST_TIMEOUT: Duration = Duration::from_secs(120);
38const APPROVAL_RELAY_TIMEOUT: Duration = Duration::from_secs(300);
39const INTERRUPT_GRACE: Duration = Duration::from_secs(5);
40const CODEX_PROVIDER_ENV: &str = "BAMBOO_CODEX_PROVIDER_KEY";
41const SESSION_STORE_FILE: &str = "codex-app-server-sessions.json";
42const TOKEN_FILE: &str = "codex-app-server-provider-token";
43const MAX_LOGICAL_SESSIONS: usize = 256;
44
45const ENV_ALLOWLIST: &[&str] = &[
46    "HOME", "PATH", "SHELL", "TERM", "LANG", "TMPDIR", "USER", "LOGNAME",
47];
48
49#[derive(Debug, Clone, Serialize, Deserialize)]
50struct AppServerSessionState {
51    thread_id: String,
52    workspace: Option<String>,
53    codex_home_mode: String,
54    updated_at: DateTime<Utc>,
55}
56
57#[derive(Debug, Default, Serialize, Deserialize)]
58struct AppServerSessionStore {
59    #[serde(default)]
60    sessions: HashMap<String, AppServerSessionState>,
61}
62
63struct AppServerConnection {
64    child: Child,
65    write_tx: mpsc::UnboundedSender<Value>,
66    incoming_rx: mpsc::Receiver<Value>,
67    pending: Arc<Mutex<HashMap<u64, oneshot::Sender<Value>>>>,
68    next_id: AtomicU64,
69    stderr_tail: Arc<Mutex<String>>,
70    writer_task: tokio::task::JoinHandle<()>,
71    reader_task: tokio::task::JoinHandle<()>,
72}
73
74impl AppServerConnection {
75    fn is_alive(&mut self) -> bool {
76        matches!(self.child.try_wait(), Ok(None))
77            && !self.writer_task.is_finished()
78            && !self.reader_task.is_finished()
79    }
80
81    fn send(&self, value: Value) -> Result<(), String> {
82        self.write_tx
83            .send(value)
84            .map_err(|_| "Codex app-server stdin writer closed".to_string())
85    }
86
87    async fn request(&self, method: &str, params: Value) -> Result<Value, String> {
88        self.request_with_timeout(method, params, REQUEST_TIMEOUT)
89            .await
90    }
91
92    async fn request_with_timeout(
93        &self,
94        method: &str,
95        params: Value,
96        timeout: Duration,
97    ) -> Result<Value, String> {
98        let id = self.next_id.fetch_add(1, Ordering::Relaxed);
99        let (reply_tx, reply_rx) = oneshot::channel();
100        self.pending.lock().await.insert(id, reply_tx);
101        if let Err(error) = self.send(json!({"id": id, "method": method, "params": params})) {
102            self.pending.lock().await.remove(&id);
103            return Err(error);
104        }
105        let response = match tokio::time::timeout(timeout, reply_rx).await {
106            Ok(Ok(response)) => response,
107            Ok(Err(_)) => {
108                return Err(format!(
109                    "Codex app-server closed while waiting for {method} response"
110                ))
111            }
112            Err(_) => {
113                self.pending.lock().await.remove(&id);
114                return Err(format!(
115                    "Codex app-server {method} request timed out after {}s",
116                    timeout.as_secs()
117                ));
118            }
119        };
120        if let Some(error) = response.get("error") {
121            return Err(format!(
122                "Codex app-server {method} failed: {}",
123                value_text(error)
124            ));
125        }
126        Ok(response.get("result").cloned().unwrap_or(Value::Null))
127    }
128
129    async fn stderr_summary(&self) -> String {
130        self.stderr_tail.lock().await.trim().to_string()
131    }
132}
133
134impl Drop for AppServerConnection {
135    fn drop(&mut self) {
136        self.writer_task.abort();
137        self.reader_task.abort();
138        #[cfg(unix)]
139        if let Some(pgid) = self.child.id().map(|pid| pid as libc::pid_t) {
140            // The executor normally uses the graceful interrupt plus
141            // TERM/KILL ladder. Drop is the last-resort worker-shutdown path,
142            // so synchronously kill the whole process group rather than
143            // allowing an in-flight Codex tool descendant to outlive Bamboo.
144            // SAFETY: a negative live child pid targets only its process group.
145            let _ = unsafe { libc::kill(-pgid, libc::SIGKILL) };
146        }
147        let _ = self.child.start_kill();
148    }
149}
150
151#[derive(Clone, Copy)]
152struct AppServerRunPolicy<'a> {
153    sandbox: &'a str,
154    approval_policy: &'a str,
155    network_access: bool,
156}
157
158#[derive(Clone, Copy)]
159enum UnexpectedApprovalDisposition {
160    Relay,
161    Deny,
162}
163
164/// `codex app-server` implementation of the Codex executor mode.
165pub struct CodexAppServerExecutor {
166    binary: PathBuf,
167    version: String,
168    model: Option<String>,
169    permissions: CodexPermissionConfig,
170    workspace: Option<String>,
171    state_dir: PathBuf,
172    forward_env: Vec<String>,
173    auth: CodexAuthConfig,
174    approval_timeout: Duration,
175    run_lock: Mutex<()>,
176    connection: Mutex<Option<AppServerConnection>>,
177}
178
179struct RunTokenGuard {
180    file: std::fs::File,
181}
182
183impl Drop for RunTokenGuard {
184    fn drop(&mut self) {
185        if let Err(error) = self.file.set_len(0) {
186            tracing::warn!(%error, "codex app-server: clear per-run provider token");
187        }
188    }
189}
190
191impl CodexAppServerExecutor {
192    #[allow(clippy::too_many_arguments)]
193    pub async fn new(
194        binary: Option<String>,
195        model: Option<String>,
196        workspace: Option<String>,
197        state_dir: Option<PathBuf>,
198        forward_env: Vec<String>,
199        auth: CodexAuthConfig,
200        permissions: CodexPermissionConfig,
201    ) -> Result<Self, String> {
202        let discovery = discover_codex_app_server(binary.as_deref()).await?;
203        let state_dir = state_dir.ok_or_else(|| {
204            "Codex app-server mode requires a Bamboo-managed state directory".to_string()
205        })?;
206        secure_directory(&state_dir).await?;
207        let executor = Self {
208            binary: PathBuf::from(discovery.path),
209            version: discovery.version,
210            model,
211            permissions,
212            workspace,
213            state_dir,
214            forward_env,
215            auth,
216            approval_timeout: APPROVAL_RELAY_TIMEOUT,
217            run_lock: Mutex::new(()),
218            connection: Mutex::new(None),
219        };
220        executor.prepare_auth_home().await?;
221        Ok(executor)
222    }
223
224    fn codex_home(&self) -> Option<PathBuf> {
225        self.auth
226            .isolated()
227            .then(|| self.state_dir.join("codex-app-server-home"))
228    }
229
230    fn codex_home_mode(&self) -> &'static str {
231        if self.auth.isolated() {
232            "isolated"
233        } else {
234            "inherit"
235        }
236    }
237
238    fn token_path(&self) -> PathBuf {
239        self.state_dir.join(TOKEN_FILE)
240    }
241
242    fn session_store_path(&self) -> PathBuf {
243        self.state_dir.join(SESSION_STORE_FILE)
244    }
245
246    async fn prepare_auth_home(&self) -> Result<(), String> {
247        let Some(home) = self.codex_home() else {
248            return Ok(());
249        };
250        secure_directory(&home).await?;
251        let auth_path = home.join("auth.json");
252        match tokio::fs::remove_file(&auth_path).await {
253            Ok(()) => {}
254            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
255            Err(error) => {
256                return Err(format!(
257                    "remove stale isolated Codex auth '{}': {error}",
258                    auth_path.display()
259                ))
260            }
261        }
262        let helper = std::env::current_exe()
263            .map_err(|error| format!("resolve Bamboo token helper executable: {error}"))?;
264        let config = self
265            .auth
266            .generated_app_server_config_toml(&helper, &self.token_path())?;
267        let config_path = home.join("config.toml");
268        tokio::fs::write(&config_path, config)
269            .await
270            .map_err(|error| {
271                format!(
272                    "write isolated Codex app-server config '{}': {error}",
273                    config_path.display()
274                )
275            })?;
276        secure_file(&config_path).await?;
277        Ok(())
278    }
279
280    fn install_run_token(&self, spec: &RunSpec) -> Result<Option<RunTokenGuard>, String> {
281        if self.auth.mode() != CodexAuthMode::Bamboo {
282            return Ok(None);
283        }
284        let token = spec
285            .secrets
286            .codex_provider_token
287            .as_ref()
288            .map(bamboo_subagent::proto::SecretValue::expose)
289            .ok_or_else(|| {
290                "Codex bamboo auth requires a fresh per-run provider token".to_string()
291            })?;
292        let mut guard = RunTokenGuard {
293            file: open_secret_for_replace(&self.token_path())?,
294        };
295        guard
296            .file
297            .write_all(token.as_bytes())
298            .map_err(|error| format!("write Codex app-server provider token: {error}"))?;
299        guard
300            .file
301            .sync_data()
302            .map_err(|error| format!("sync Codex app-server provider token: {error}"))?;
303        Ok(Some(guard))
304    }
305
306    async fn load_session_store(&self) -> AppServerSessionStore {
307        let Ok(bytes) = tokio::fs::read(self.session_store_path()).await else {
308            return AppServerSessionStore::default();
309        };
310        serde_json::from_slice(&bytes).unwrap_or_default()
311    }
312
313    async fn save_session_store(&self, store: &AppServerSessionStore) {
314        if let Err(error) = write_json_atomic(&self.session_store_path(), store).await {
315            tracing::warn!(%error, "codex app-server: persist logical session map");
316        }
317    }
318
319    async fn stored_thread(&self, logical_session: &str) -> Option<String> {
320        let store = self.load_session_store().await;
321        let state = store.sessions.get(logical_session)?;
322        if state.workspace != self.workspace || state.codex_home_mode != self.codex_home_mode() {
323            return None;
324        }
325        (!state.thread_id.trim().is_empty()).then(|| state.thread_id.clone())
326    }
327
328    async fn store_thread(&self, logical_session: &str, thread_id: &str) {
329        let mut store = self.load_session_store().await;
330        store.sessions.insert(
331            logical_session.to_string(),
332            AppServerSessionState {
333                thread_id: thread_id.to_string(),
334                workspace: self.workspace.clone(),
335                codex_home_mode: self.codex_home_mode().to_string(),
336                updated_at: Utc::now(),
337            },
338        );
339        prune_session_store(&mut store, logical_session);
340        self.save_session_store(&store).await;
341    }
342
343    async fn forget_thread(&self, logical_session: &str) {
344        let mut store = self.load_session_store().await;
345        if store.sessions.remove(logical_session).is_some() {
346            self.save_session_store(&store).await;
347        }
348    }
349
350    fn build_command(&self) -> Result<Command, String> {
351        let mut command = Command::new(&self.binary);
352        command.arg("app-server").arg("--listen").arg("stdio://");
353        command.env_clear();
354        for (key, value) in std::env::vars() {
355            if ENV_ALLOWLIST.contains(&key.as_str()) || key.starts_with("LC_") {
356                command.env(key, value);
357            }
358        }
359        if let Some(home) = self.codex_home() {
360            command.env("CODEX_HOME", home);
361        }
362        for name in &self.forward_env {
363            if let Ok(value) = std::env::var(name) {
364                command.env(name, value);
365            }
366        }
367        if self.auth.mode() == CodexAuthMode::Custom {
368            let key = self.auth.provider_key().ok_or_else(|| {
369                "Codex custom provider key was not resolved at provisioning".to_string()
370            })?;
371            command.env(CODEX_PROVIDER_ENV, key);
372        }
373        command.stdin(Stdio::piped());
374        command.stdout(Stdio::piped());
375        command.stderr(Stdio::piped());
376        command.kill_on_drop(true);
377        #[cfg(unix)]
378        command.process_group(0);
379        Ok(command)
380    }
381
382    async fn start_connection(&self) -> Result<AppServerConnection, String> {
383        self.prepare_auth_home().await?;
384        let mut child = self.build_command()?.spawn().map_err(|error| {
385            format!(
386                "spawn Codex app-server '{}': {error}",
387                self.binary.display()
388            )
389        })?;
390        let stdin = child
391            .stdin
392            .take()
393            .ok_or_else(|| "Codex app-server has no stdin pipe".to_string())?;
394        let stdout = child
395            .stdout
396            .take()
397            .ok_or_else(|| "Codex app-server has no stdout pipe".to_string())?;
398        let stderr = child.stderr.take();
399
400        let (write_tx, mut write_rx) = mpsc::unbounded_channel::<Value>();
401        let writer_task = tokio::spawn(async move {
402            let mut stdin = stdin;
403            while let Some(value) = write_rx.recv().await {
404                let Ok(mut bytes) = serde_json::to_vec(&value) else {
405                    continue;
406                };
407                bytes.push(b'\n');
408                if stdin.write_all(&bytes).await.is_err() || stdin.flush().await.is_err() {
409                    break;
410                }
411            }
412        });
413
414        let pending: Arc<Mutex<HashMap<u64, oneshot::Sender<Value>>>> =
415            Arc::new(Mutex::new(HashMap::new()));
416        // Match the proven cc-connect posture: bounded event buffering applies
417        // backpressure instead of allowing a noisy app-server to grow memory
418        // without limit. Client responses bypass this queue via `pending`.
419        let (incoming_tx, incoming_rx) = mpsc::channel(128);
420        let reader_pending = pending.clone();
421        let reader_task = tokio::spawn(async move {
422            let mut reader = BufReader::new(stdout);
423            loop {
424                let line = match read_bounded_line(&mut reader, MAX_STDOUT_LINE_BYTES).await {
425                    Ok(Some(line)) => line,
426                    Ok(None) => break,
427                    Err(error) => {
428                        let _ = incoming_tx
429                            .send(json!({
430                                "method": "bamboo/transport/error",
431                                "params": {"message": error.to_string()}
432                            }))
433                            .await;
434                        break;
435                    }
436                };
437                let value: Value = match serde_json::from_slice(&line) {
438                    Ok(value) => value,
439                    Err(error) => {
440                        let _ = incoming_tx
441                            .send(json!({
442                                "method": "bamboo/transport/error",
443                                "params": {"message": format!("invalid JSONL: {error}")}
444                            }))
445                            .await;
446                        break;
447                    }
448                };
449                if value.get("method").is_none() {
450                    if let Some(id) = value.get("id").and_then(Value::as_u64) {
451                        if let Some(reply) = reader_pending.lock().await.remove(&id) {
452                            let _ = reply.send(value);
453                            continue;
454                        }
455                    }
456                }
457                if incoming_tx.send(value).await.is_err() {
458                    break;
459                }
460            }
461            reader_pending.lock().await.clear();
462        });
463
464        let stderr_tail = Arc::new(Mutex::new(String::new()));
465        if let Some(stderr) = stderr {
466            let tail = stderr_tail.clone();
467            tokio::spawn(async move { drain_stderr_tail(stderr, tail).await });
468        }
469        let connection = AppServerConnection {
470            child,
471            write_tx,
472            incoming_rx,
473            pending,
474            next_id: AtomicU64::new(1),
475            stderr_tail,
476            writer_task,
477            reader_task,
478        };
479        connection
480            .request(
481                "initialize",
482                json!({
483                    "clientInfo": {
484                        "name": "bamboo",
485                        "title": "Bamboo",
486                        "version": env!("CARGO_PKG_VERSION")
487                    },
488                    "capabilities": {"experimentalApi": true}
489                }),
490            )
491            .await?;
492        connection.send(json!({"method": "initialized", "params": {}}))?;
493        Ok(connection)
494    }
495
496    async fn ensure_connection<'a>(
497        &'a self,
498        slot: &'a mut Option<AppServerConnection>,
499    ) -> Result<&'a mut AppServerConnection, String> {
500        let alive = slot.as_mut().is_some_and(AppServerConnection::is_alive);
501        if !alive {
502            if let Some(mut stale) = slot.take() {
503                terminate_child(&mut stale.child).await;
504            }
505            *slot = Some(self.start_connection().await?);
506        }
507        Ok(slot.as_mut().expect("connection installed"))
508    }
509
510    fn thread_params(&self, policy: AppServerRunPolicy<'_>) -> Value {
511        let mut params = json!({
512            "approvalPolicy": policy.approval_policy,
513            "approvalsReviewer": (policy.approval_policy != "never").then_some("user"),
514            "sandbox": policy.sandbox,
515            "cwd": self.workspace,
516            "model": self.model,
517        });
518        remove_null_object_fields(&mut params);
519        params
520    }
521
522    async fn start_thread(
523        &self,
524        connection: &AppServerConnection,
525        policy: AppServerRunPolicy<'_>,
526    ) -> Result<String, String> {
527        let result = connection
528            .request("thread/start", self.thread_params(policy))
529            .await?;
530        result
531            .pointer("/thread/id")
532            .and_then(Value::as_str)
533            .filter(|id| !id.is_empty())
534            .map(str::to_string)
535            .ok_or_else(|| "Codex thread/start response omitted thread.id".to_string())
536    }
537
538    async fn resume_thread(
539        &self,
540        connection: &AppServerConnection,
541        thread_id: &str,
542        policy: AppServerRunPolicy<'_>,
543    ) -> Result<(), String> {
544        let mut params = self.thread_params(policy);
545        params["threadId"] = Value::String(thread_id.to_string());
546        connection
547            .request("thread/resume", params)
548            .await
549            .map(|_| ())
550    }
551
552    async fn start_turn(
553        &self,
554        connection: &AppServerConnection,
555        thread_id: &str,
556        prompt: &str,
557        reasoning_effort: Option<&str>,
558        policy: AppServerRunPolicy<'_>,
559    ) -> Result<String, String> {
560        let mut params = json!({
561            "threadId": thread_id,
562            "input": [{"type": "text", "text": prompt, "text_elements": []}],
563            "approvalPolicy": policy.approval_policy,
564            "approvalsReviewer": (policy.approval_policy != "never").then_some("user"),
565            "cwd": self.workspace,
566            "model": self.model,
567            "effort": reasoning_effort,
568            "sandboxPolicy": sandbox_policy(
569                policy.sandbox,
570                self.workspace.as_deref(),
571                policy.network_access,
572            ),
573        });
574        remove_null_object_fields(&mut params);
575        let result = connection.request("turn/start", params).await?;
576        result
577            .pointer("/turn/id")
578            .and_then(Value::as_str)
579            .filter(|id| !id.is_empty())
580            .map(str::to_string)
581            .ok_or_else(|| "Codex turn/start response omitted turn.id".to_string())
582    }
583
584    #[allow(clippy::too_many_arguments)]
585    async fn drive_turn(
586        &self,
587        connection: &mut AppServerConnection,
588        thread_id: &str,
589        turn_id: &str,
590        events: &EventSink,
591        steer: &mut SteerInbox,
592        approval_tasks: &mut Vec<tokio::task::JoinHandle<()>>,
593        force_cancelled: bool,
594        approval_disposition: UnexpectedApprovalDisposition,
595    ) -> ChildOutcome {
596        let mut state = AppRunState::default();
597        let mut steer_open = true;
598        loop {
599            tokio::select! {
600                maybe_message = steer.recv(), if steer_open => {
601                    if let Some(message) = maybe_message {
602                        if let Err(error) = connection.request("turn/steer", json!({
603                            "threadId": thread_id,
604                            "expectedTurnId": turn_id,
605                            "input": [{"type": "text", "text": message, "text_elements": []}],
606                        })).await {
607                            events.emit(json!({
608                                "type": "runner_progress",
609                                "session_id": thread_id,
610                                "round_count": 1,
611                                "executor": "codex_app_server",
612                                "phase": "steer_rejected",
613                                "message": error,
614                            }));
615                        }
616                    } else {
617                        steer_open = false;
618                    }
619                }
620                incoming = connection.incoming_rx.recv() => {
621                    let Some(message) = incoming else {
622                        let stderr = connection.stderr_summary().await;
623                        let suffix = if stderr.is_empty() { String::new() } else { format!("; stderr: {stderr}") };
624                        return ChildOutcome::error(format!("Codex app-server transport closed{suffix}"));
625                    };
626                    let method = message.get("method").and_then(Value::as_str).unwrap_or("");
627                    if message.get("id").is_some() {
628                        if is_approval_method(method) {
629                            if !approval_matches_active_turn(&message, thread_id, turn_id) {
630                                let _ = connection.send(approval_response(&message, false));
631                                continue;
632                            }
633                            if method == "item/permissions/requestApproval" {
634                                // Granular filesystem/network escalation cannot be
635                                // represented by Bamboo's boolean relay without
636                                // accidentally granting the full requested set.
637                                // An empty subset is the protocol's deterministic
638                                // deny response in every posture.
639                                let _ = connection.send(approval_response(&message, false));
640                                continue;
641                            }
642                            match approval_disposition {
643                                UnexpectedApprovalDisposition::Deny => {
644                                    let _ = connection.send(approval_response(&message, false));
645                                }
646                                UnexpectedApprovalDisposition::Relay => {
647                                    approval_tasks.push(spawn_approval_relay(
648                                        connection.write_tx.clone(),
649                                        message,
650                                        events.host().cloned(),
651                                        self.approval_timeout,
652                                    ));
653                                }
654                            }
655                        } else {
656                            let id = message.get("id").cloned().unwrap_or(Value::Null);
657                            let _ = connection.send(json!({
658                                "id": id,
659                                "error": {"code": -32601, "message": format!("unsupported server request: {method}")}
660                            }));
661                        }
662                        continue;
663                    }
664                    if method == "bamboo/transport/error" {
665                        return ChildOutcome::error(error_message(&message, "Codex app-server transport error"));
666                    }
667                    if let Some(outcome) = handle_notification(
668                        method,
669                        message.get("params").unwrap_or(&Value::Null),
670                        thread_id,
671                        turn_id,
672                        events,
673                        &mut state,
674                        force_cancelled,
675                    ) {
676                        return outcome;
677                    }
678                }
679            }
680        }
681    }
682
683    async fn run_inner(
684        &self,
685        spec: &RunSpec,
686        events: &EventSink,
687        steer: &mut SteerInbox,
688        cancel: &CancellationToken,
689    ) -> ChildOutcome {
690        let logical_session = match spec
691            .permission_policy
692            .as_ref()
693            .map(|policy| policy.session_id.trim())
694            .filter(|id| !id.is_empty())
695            .map(str::to_string)
696            .or_else(|| {
697                spec.logical_session
698                    .as_ref()
699                    .map(|identity| identity.session_id.trim())
700                    .filter(|id| !id.is_empty())
701                    .map(str::to_string)
702            })
703            .or_else(|| {
704                self.permissions
705                    .provisioned_session_id()
706                    .map(str::to_string)
707            }) {
708            Some(id) => id,
709            None => {
710                return ChildOutcome::error(
711                    "Codex app-server mode requires a logical or provisioned session id",
712                )
713            }
714        };
715        let activation = match self
716            .permissions
717            .activation_permission(spec.permission_policy.as_ref())
718        {
719            Ok(activation) => activation,
720            Err(error) => {
721                events.emit(event_json(AgentEvent::Error {
722                    message: error.clone(),
723                }));
724                return ChildOutcome::error(error);
725            }
726        };
727        if let Some(reason) = activation.explicit_deny_reason.as_ref() {
728            let message = format!(
729                "Codex app-server cannot safely enforce Bamboo explicit-deny policy: {reason}"
730            );
731            events.emit(event_json(AgentEvent::PermissionPostureActivated {
732                session_id: logical_session.clone(),
733                policy_revision: activation.policy_revision,
734                requested_mode: activation.resolution.requested.as_str().to_string(),
735                effective_mode: activation.resolution.effective.as_str().to_string(),
736                executor_mapping: "codex_app_server:blocked_explicit_deny".to_string(),
737            }));
738            events.emit(event_json(AgentEvent::Error {
739                message: message.clone(),
740            }));
741            return ChildOutcome::error(message);
742        }
743        let approval_policy = if activation.resolution.suppress_approval_prompts()
744            || activation.resolution.effective == bamboo_domain::PermissionMode::Plan
745        {
746            "never"
747        } else {
748            "on-request"
749        };
750        let approval_disposition =
751            if activation.resolution.effective == bamboo_domain::PermissionMode::Plan {
752                UnexpectedApprovalDisposition::Deny
753            } else if activation.resolution.suppress_approval_prompts() {
754                // `approvalPolicy=never` means operate within the already-selected
755                // sandbox. A surprise approval may request additional filesystem or
756                // network access, so no-prompt mode must deny rather than expand it.
757                UnexpectedApprovalDisposition::Deny
758            } else {
759                UnexpectedApprovalDisposition::Relay
760            };
761        let executor_mapping = format!("codex_app_server:approvalPolicy={approval_policy}");
762        let (sandbox, network_access, warnings) =
763            self.permissions.app_server_posture(activation.resolution);
764        let run_policy = AppServerRunPolicy {
765            sandbox: &sandbox,
766            approval_policy,
767            network_access,
768        };
769        events.emit(event_json(AgentEvent::PermissionPostureActivated {
770            session_id: logical_session.clone(),
771            policy_revision: activation.policy_revision,
772            requested_mode: activation.resolution.requested.as_str().to_string(),
773            effective_mode: activation.resolution.effective.as_str().to_string(),
774            executor_mapping: executor_mapping.clone(),
775        }));
776        for warning in warnings {
777            events.emit(json!({
778                "type": "runner_progress",
779                "session_id": logical_session,
780                "round_count": 0,
781                "level": "warning",
782                "message": warning,
783            }));
784        }
785        events.emit(json!({
786            "type": "runner_progress",
787            "session_id": logical_session,
788            "round_count": 0,
789            "executor": "codex_app_server",
790            "binary": self.binary,
791            "version": self.version,
792            "model": self.model,
793            "auth_mode": self.auth.mode().as_str(),
794            "codex_home_mode": self.codex_home_mode(),
795            "sandbox": sandbox,
796            "approval_policy": approval_policy,
797            "approvals_reviewer": (!activation.resolution.suppress_approval_prompts()
798                && activation.resolution.effective != bamboo_domain::PermissionMode::Plan)
799                .then_some("user"),
800            "network_access": network_access,
801            "permission_profile": self.permissions.permission_profile(),
802            "requested_mode": activation.resolution.requested.as_str(),
803            "effective_mode": activation.resolution.effective.as_str(),
804            "executor_mapping": executor_mapping,
805        }));
806
807        if spec.messages.is_empty() {
808            self.forget_thread(&logical_session).await;
809        }
810        let mut slot = self.connection.lock().await;
811        let connection = match self.ensure_connection(&mut slot).await {
812            Ok(connection) => connection,
813            Err(error) => return ChildOutcome::error(error),
814        };
815
816        while connection.incoming_rx.try_recv().is_ok() {}
817        let stored = if spec.messages.is_empty() {
818            None
819        } else {
820            self.stored_thread(&logical_session).await
821        };
822        let (thread_id, prompt) = if let Some(thread_id) = stored {
823            match self.resume_thread(connection, &thread_id, run_policy).await {
824                Ok(()) => (thread_id, spec.assignment.clone()),
825                Err(error) => {
826                    tracing::warn!(%error, "codex app-server: resume failed; rehydrating once");
827                    events.emit(json!({
828                        "type": "runner_progress",
829                        "session_id": logical_session,
830                        "round_count": 0,
831                        "executor": "codex_app_server",
832                        "phase": "resume_fallback",
833                        "message": "resume failed; starting a new thread with bounded history rehydration",
834                    }));
835                    self.forget_thread(&logical_session).await;
836                    let new_id = match self.start_thread(connection, run_policy).await {
837                        Ok(id) => id,
838                        Err(error) => return ChildOutcome::error(error),
839                    };
840                    (
841                        new_id,
842                        build_rehydrated_turn(&spec.messages, &spec.assignment),
843                    )
844                }
845            }
846        } else {
847            let thread_id = match self.start_thread(connection, run_policy).await {
848                Ok(id) => id,
849                Err(error) => return ChildOutcome::error(error),
850            };
851            let prompt = if spec.messages.is_empty() {
852                spec.assignment.clone()
853            } else {
854                build_rehydrated_turn(&spec.messages, &spec.assignment)
855            };
856            (thread_id, prompt)
857        };
858        self.store_thread(&logical_session, &thread_id).await;
859
860        let turn_id = match self
861            .start_turn(
862                connection,
863                &thread_id,
864                &prompt,
865                spec.reasoning_effort.as_deref(),
866                run_policy,
867            )
868            .await
869        {
870            Ok(id) => id,
871            Err(error) => return ChildOutcome::error(error),
872        };
873        events.emit(event_json(AgentEvent::RunnerProgress {
874            session_id: thread_id.clone(),
875            round_count: 1,
876        }));
877        let mut approval_tasks = Vec::new();
878        let outcome = tokio::select! {
879            outcome = self.drive_turn(
880                connection,
881                &thread_id,
882                &turn_id,
883                events,
884                steer,
885                &mut approval_tasks,
886                false,
887                approval_disposition,
888            ) => outcome,
889            _ = cancel.cancelled() => {
890                let interrupt = connection.request_with_timeout(
891                    "turn/interrupt",
892                    json!({"threadId": thread_id, "turnId": turn_id}),
893                    INTERRUPT_GRACE,
894                ).await;
895                if let Err(error) = interrupt {
896                    tracing::warn!(%error, "codex app-server: graceful interrupt failed");
897                }
898                match tokio::time::timeout(
899                    INTERRUPT_GRACE,
900                    self.drive_turn(
901                        connection,
902                        &thread_id,
903                        &turn_id,
904                        events,
905                        steer,
906                        &mut approval_tasks,
907                        true,
908                        approval_disposition,
909                    ),
910                ).await {
911                    Ok(_) => ChildOutcome::cancelled(),
912                    Err(_) => {
913                        if let Some(mut connection) = slot.take() {
914                            terminate_child(&mut connection.child).await;
915                        }
916                        ChildOutcome::cancelled()
917                    }
918                }
919            }
920        };
921        for task in approval_tasks {
922            task.abort();
923        }
924        outcome
925    }
926}
927
928#[async_trait]
929impl ChildExecutor for CodexAppServerExecutor {
930    async fn run(
931        &self,
932        spec: RunSpec,
933        events: EventSink,
934        mut steer: SteerInbox,
935        cancel: CancellationToken,
936    ) -> ChildOutcome {
937        // A warm executor owns one long-lived app-server and one refreshable
938        // token file. Serialize activations so a queued run cannot overwrite
939        // or clear the token still in use by the active run.
940        let _run_guard = self.run_lock.lock().await;
941        let _token_guard = match self.install_run_token(&spec) {
942            Ok(guard) => guard,
943            Err(error) => return ChildOutcome::error(error),
944        };
945        let outcome = self.run_inner(&spec, &events, &mut steer, &cancel).await;
946        outcome
947    }
948}
949
950#[derive(Default)]
951struct AppRunState {
952    last_agent_message: String,
953    last_agent_item_id: Option<String>,
954    usage: TokenUsage,
955    started_items: HashSet<String>,
956}
957
958fn handle_notification(
959    method: &str,
960    params: &Value,
961    thread_id: &str,
962    turn_id: &str,
963    events: &EventSink,
964    state: &mut AppRunState,
965    force_cancelled: bool,
966) -> Option<ChildOutcome> {
967    if params
968        .get("threadId")
969        .and_then(Value::as_str)
970        .is_some_and(|id| id != thread_id)
971        || params
972            .get("turnId")
973            .and_then(Value::as_str)
974            .is_some_and(|id| id != turn_id)
975    {
976        return None;
977    }
978    match method {
979        "item/agentMessage/delta" => {
980            if let Some(delta) = params.get("delta").and_then(Value::as_str) {
981                let item_id = params
982                    .get("itemId")
983                    .and_then(Value::as_str)
984                    .unwrap_or("codex-agent-message");
985                if state.last_agent_item_id.as_deref() != Some(item_id) {
986                    state.last_agent_item_id = Some(item_id.to_string());
987                    state.last_agent_message.clear();
988                }
989                state.last_agent_message.push_str(delta);
990                events.emit(event_json(AgentEvent::Token {
991                    content: delta.to_string(),
992                }));
993            }
994        }
995        "item/reasoning/textDelta" | "item/reasoning/summaryTextDelta" => {
996            if let Some(delta) = params.get("delta").and_then(Value::as_str) {
997                events.emit(event_json(AgentEvent::ReasoningToken {
998                    content: delta.to_string(),
999                }));
1000            }
1001        }
1002        "item/commandExecution/outputDelta" | "item/fileChange/outputDelta" => {
1003            if let (Some(item_id), Some(delta)) = (
1004                params.get("itemId").and_then(Value::as_str),
1005                params.get("delta").and_then(Value::as_str),
1006            ) {
1007                events.emit(event_json(AgentEvent::ToolToken {
1008                    tool_call_id: item_id.to_string(),
1009                    content: delta.to_string(),
1010                }));
1011            }
1012        }
1013        "item/mcpToolCall/progress" => {
1014            if let (Some(item_id), Some(message)) = (
1015                params.get("itemId").and_then(Value::as_str),
1016                params.get("message").and_then(Value::as_str),
1017            ) {
1018                events.emit(event_json(AgentEvent::ToolToken {
1019                    tool_call_id: item_id.to_string(),
1020                    content: message.to_string(),
1021                }));
1022            }
1023        }
1024        "item/started" => {
1025            if let Some(item) = params.get("item") {
1026                emit_item_started(item, events, state);
1027            }
1028        }
1029        "item/completed" => {
1030            if let Some(item) = params.get("item") {
1031                emit_item_completed(item, events, state);
1032            }
1033        }
1034        "thread/tokenUsage/updated" => {
1035            state.usage = parse_app_server_usage(params.get("tokenUsage"));
1036        }
1037        "turn/completed" => {
1038            if force_cancelled {
1039                events.emit(event_json(AgentEvent::Cancelled {
1040                    message: Some("Codex app-server turn interrupted".to_string()),
1041                }));
1042                return Some(ChildOutcome::cancelled());
1043            }
1044            let turn = params.get("turn").unwrap_or(&Value::Null);
1045            match turn.get("status").and_then(Value::as_str) {
1046                Some("completed") => {
1047                    events.emit(event_json(AgentEvent::Complete { usage: state.usage }));
1048                    return Some(ChildOutcome::completed(state.last_agent_message.clone()));
1049                }
1050                Some("interrupted") => {
1051                    events.emit(event_json(AgentEvent::Cancelled {
1052                        message: Some("Codex app-server turn interrupted".to_string()),
1053                    }));
1054                    return Some(ChildOutcome::cancelled());
1055                }
1056                _ => {
1057                    let message = error_message(turn, "Codex app-server turn failed");
1058                    events.emit(event_json(AgentEvent::Error {
1059                        message: message.clone(),
1060                    }));
1061                    return Some(ChildOutcome::error(message));
1062                }
1063            }
1064        }
1065        "error" => {
1066            let message = error_message(params, "Codex app-server error");
1067            events.emit(event_json(AgentEvent::Error {
1068                message: message.clone(),
1069            }));
1070            return Some(ChildOutcome::error(message));
1071        }
1072        other => tracing::debug!(
1073            method = other,
1074            "codex app-server: unrecognized notification"
1075        ),
1076    }
1077    None
1078}
1079
1080fn emit_item_started(item: &Value, events: &EventSink, state: &mut AppRunState) {
1081    let item_id = item
1082        .get("id")
1083        .and_then(Value::as_str)
1084        .unwrap_or("codex-item");
1085    if !state.started_items.insert(item_id.to_string()) {
1086        return;
1087    }
1088    let (tool_name, arguments) = match item.get("type").and_then(Value::as_str) {
1089        Some("commandExecution") => (
1090            "Bash".to_string(),
1091            json!({"command": item.get("command"), "cwd": item.get("cwd")}),
1092        ),
1093        Some("fileChange") => (
1094            "ApplyPatch".to_string(),
1095            json!({"changes": item.get("changes")}),
1096        ),
1097        Some("mcpToolCall") => (
1098            format!(
1099                "{}::{}",
1100                item.get("server").and_then(Value::as_str).unwrap_or("mcp"),
1101                item.get("tool").and_then(Value::as_str).unwrap_or("tool")
1102            ),
1103            item.get("arguments").cloned().unwrap_or_else(|| json!({})),
1104        ),
1105        Some("dynamicToolCall") => (
1106            item.get("tool")
1107                .or_else(|| item.get("name"))
1108                .and_then(Value::as_str)
1109                .unwrap_or("DynamicTool")
1110                .to_string(),
1111            item.get("arguments").cloned().unwrap_or_else(|| json!({})),
1112        ),
1113        Some("webSearch") => ("WebSearch".to_string(), json!({"query": item.get("query")})),
1114        _ => return,
1115    };
1116    events.emit(event_json(AgentEvent::ToolStart {
1117        tool_call_id: item_id.to_string(),
1118        tool_name,
1119        arguments,
1120    }));
1121}
1122
1123fn emit_item_completed(item: &Value, events: &EventSink, state: &mut AppRunState) {
1124    if item.get("type").and_then(Value::as_str) == Some("agentMessage") {
1125        if let Some(text) = item.get("text").and_then(Value::as_str) {
1126            state.last_agent_item_id = item.get("id").and_then(Value::as_str).map(str::to_string);
1127            state.last_agent_message = text.to_string();
1128        }
1129        return;
1130    }
1131    emit_item_started(item, events, state);
1132    let item_id = item
1133        .get("id")
1134        .and_then(Value::as_str)
1135        .unwrap_or("codex-item");
1136    if !state.started_items.contains(item_id) {
1137        return;
1138    }
1139    let status = item.get("status").and_then(Value::as_str).unwrap_or("");
1140    let error = item
1141        .get("error")
1142        .filter(|value| !value.is_null())
1143        .map(value_text)
1144        .or_else(|| {
1145            matches!(status, "failed" | "declined" | "error").then(|| {
1146                item.get("aggregatedOutput")
1147                    .map(value_text)
1148                    .unwrap_or_else(|| format!("Codex tool finished with status {status}"))
1149            })
1150        });
1151    if let Some(error) = error {
1152        events.emit(event_json(AgentEvent::ToolError {
1153            tool_call_id: item_id.to_string(),
1154            error: truncate_chars(&error, TOOL_RESULT_TRUNCATE_CHARS),
1155        }));
1156    } else {
1157        let result = item
1158            .get("aggregatedOutput")
1159            .or_else(|| item.get("result"))
1160            .or_else(|| item.get("changes"))
1161            .map(value_text)
1162            .unwrap_or_else(|| status.to_string());
1163        events.emit(event_json(AgentEvent::ToolComplete {
1164            tool_call_id: item_id.to_string(),
1165            result: ToolResult::text(true, truncate_chars(&result, TOOL_RESULT_TRUNCATE_CHARS)),
1166        }));
1167    }
1168}
1169
1170fn is_approval_method(method: &str) -> bool {
1171    matches!(
1172        method,
1173        "item/commandExecution/requestApproval"
1174            | "item/fileChange/requestApproval"
1175            | "item/permissions/requestApproval"
1176            | "execCommandApproval"
1177            | "applyPatchApproval"
1178    )
1179}
1180
1181fn spawn_approval_relay(
1182    write_tx: mpsc::UnboundedSender<Value>,
1183    request: Value,
1184    host: Option<HostBridge>,
1185    timeout: Duration,
1186) -> tokio::task::JoinHandle<()> {
1187    tokio::spawn(async move {
1188        let id = request.get("id").cloned().unwrap_or(Value::Null);
1189        let method = request
1190            .get("method")
1191            .and_then(Value::as_str)
1192            .unwrap_or("")
1193            .to_string();
1194        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1195        let is_command = method.contains("commandExecution") || method == "execCommandApproval";
1196        let body = if is_command {
1197            json!({
1198                "tool_name": "Bash",
1199                "permission_type": "command_execution",
1200                "resource": params.get("command").cloned().unwrap_or(Value::Null),
1201                "question": params.get("reason").cloned().unwrap_or_else(|| Value::String("Codex requests permission to execute a command".to_string())),
1202                "input": params,
1203            })
1204        } else {
1205            json!({
1206                "tool_name": "ApplyPatch",
1207                "permission_type": "file_change",
1208                "resource": params.get("grantRoot").or_else(|| params.get("path")).cloned().unwrap_or(Value::Null),
1209                "question": params.get("reason").cloned().unwrap_or_else(|| Value::String("Codex requests permission to modify files".to_string())),
1210                "input": params,
1211            })
1212        };
1213        let approved = if let Some(host) = host {
1214            match tokio::time::timeout(timeout, host.approval_call(body)).await {
1215                Ok(Ok(reply)) => reply
1216                    .get("approved")
1217                    .and_then(Value::as_bool)
1218                    .unwrap_or(false),
1219                Ok(Err(error)) => {
1220                    tracing::warn!(%error, "codex app-server: approval relay failed closed");
1221                    false
1222                }
1223                Err(_) => {
1224                    tracing::warn!(
1225                        seconds = timeout.as_secs(),
1226                        "codex app-server: approval relay timed out; denying"
1227                    );
1228                    false
1229                }
1230            }
1231        } else {
1232            tracing::warn!("codex app-server: approval host bridge unavailable; denying");
1233            false
1234        };
1235        let _ = write_tx.send(approval_response_parts(id, &method, approved));
1236    })
1237}
1238
1239fn approval_matches_active_turn(request: &Value, thread_id: &str, turn_id: &str) -> bool {
1240    let params = request.get("params").unwrap_or(&Value::Null);
1241    !params
1242        .get("threadId")
1243        .and_then(Value::as_str)
1244        .is_some_and(|id| id != thread_id)
1245        && !params
1246            .get("turnId")
1247            .and_then(Value::as_str)
1248            .is_some_and(|id| id != turn_id)
1249}
1250
1251fn approval_response(request: &Value, approved: bool) -> Value {
1252    approval_response_parts(
1253        request.get("id").cloned().unwrap_or(Value::Null),
1254        request.get("method").and_then(Value::as_str).unwrap_or(""),
1255        approved,
1256    )
1257}
1258
1259fn approval_response_parts(id: Value, method: &str, approved: bool) -> Value {
1260    if method == "item/permissions/requestApproval" {
1261        // Omitted granular permissions are denied by the app-server protocol.
1262        // Bamboo deliberately grants an empty subset because its approval bridge
1263        // cannot safely express per-item FS/network escalation.
1264        return json!({"id": id, "result": {"permissions": {}}});
1265    }
1266    let decision = if matches!(method, "execCommandApproval" | "applyPatchApproval") {
1267        if approved {
1268            "approved"
1269        } else {
1270            "denied"
1271        }
1272    } else if approved {
1273        "accept"
1274    } else {
1275        "decline"
1276    };
1277    json!({"id": id, "result": {"decision": decision}})
1278}
1279
1280fn sandbox_policy(sandbox: &str, workspace: Option<&str>, network_access: bool) -> Value {
1281    match sandbox {
1282        "read-only" => json!({"type": "readOnly", "networkAccess": false}),
1283        "danger-full-access" => json!({"type": "dangerFullAccess"}),
1284        _ => json!({
1285            "type": "workspaceWrite",
1286            "writableRoots": workspace.into_iter().collect::<Vec<_>>(),
1287            "networkAccess": network_access,
1288        }),
1289    }
1290}
1291
1292fn parse_app_server_usage(value: Option<&Value>) -> TokenUsage {
1293    let usage = value
1294        .and_then(|value| value.get("last"))
1295        .unwrap_or(&Value::Null);
1296    TokenUsage {
1297        prompt_tokens: usage
1298            .get("inputTokens")
1299            .and_then(Value::as_u64)
1300            .unwrap_or(0),
1301        completion_tokens: usage
1302            .get("outputTokens")
1303            .and_then(Value::as_u64)
1304            .unwrap_or(0),
1305        total_tokens: usage
1306            .get("totalTokens")
1307            .and_then(Value::as_u64)
1308            .unwrap_or(0),
1309    }
1310}
1311
1312fn event_json(event: AgentEvent) -> Value {
1313    serde_json::to_value(event)
1314        .unwrap_or_else(|_| json!({"type": "error", "message": "serialize agent event"}))
1315}
1316
1317fn error_message(value: &Value, fallback: &str) -> String {
1318    value
1319        .pointer("/error/message")
1320        .or_else(|| value.get("message"))
1321        .or_else(|| value.get("error"))
1322        .map(value_text)
1323        .filter(|message| !message.is_empty())
1324        .unwrap_or_else(|| fallback.to_string())
1325}
1326
1327fn value_text(value: &Value) -> String {
1328    match value {
1329        Value::String(text) => text.clone(),
1330        Value::Null => String::new(),
1331        value => value.to_string(),
1332    }
1333}
1334
1335fn truncate_chars(value: &str, max_chars: usize) -> String {
1336    if value.chars().count() <= max_chars {
1337        return value.to_string();
1338    }
1339    let mut output: String = value.chars().take(max_chars).collect();
1340    output.push_str("\n… truncated by Bamboo …");
1341    output
1342}
1343
1344fn remove_null_object_fields(value: &mut Value) {
1345    if let Some(object) = value.as_object_mut() {
1346        object.retain(|_, value| !value.is_null());
1347    }
1348}
1349
1350fn prune_session_store(store: &mut AppServerSessionStore, current_session: &str) {
1351    while store.sessions.len() > MAX_LOGICAL_SESSIONS {
1352        let oldest = store
1353            .sessions
1354            .iter()
1355            .filter(|(session, _)| session.as_str() != current_session)
1356            .min_by_key(|(_, state)| state.updated_at)
1357            .map(|(session, _)| session.clone());
1358        let Some(oldest) = oldest else {
1359            break;
1360        };
1361        store.sessions.remove(&oldest);
1362    }
1363}
1364
1365fn open_secret_for_replace(path: &Path) -> Result<std::fs::File, String> {
1366    let mut options = std::fs::OpenOptions::new();
1367    options.create(true).truncate(true).write(true);
1368    #[cfg(unix)]
1369    {
1370        use std::os::unix::fs::OpenOptionsExt as _;
1371        options
1372            .mode(0o600)
1373            .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC);
1374    }
1375    let file = options
1376        .open(path)
1377        .map_err(|error| format!("open Codex app-server provider token: {error}"))?;
1378    if !file
1379        .metadata()
1380        .map_err(|error| format!("inspect Codex app-server provider token: {error}"))?
1381        .is_file()
1382    {
1383        return Err("Codex app-server provider token path is not a regular file".to_string());
1384    }
1385    #[cfg(unix)]
1386    file.set_permissions(std::os::unix::fs::PermissionsExt::from_mode(0o600))
1387        .map_err(|error| format!("secure Codex app-server provider token: {error}"))?;
1388    Ok(file)
1389}
1390
1391async fn secure_directory(path: &Path) -> Result<(), String> {
1392    tokio::fs::create_dir_all(path).await.map_err(|error| {
1393        format!(
1394            "create Codex app-server state '{}': {error}",
1395            path.display()
1396        )
1397    })?;
1398    #[cfg(unix)]
1399    tokio::fs::set_permissions(path, std::os::unix::fs::PermissionsExt::from_mode(0o700))
1400        .await
1401        .map_err(|error| {
1402            format!(
1403                "secure Codex app-server state '{}': {error}",
1404                path.display()
1405            )
1406        })?;
1407    Ok(())
1408}
1409
1410async fn secure_file(path: &Path) -> Result<(), String> {
1411    #[cfg(unix)]
1412    tokio::fs::set_permissions(path, std::os::unix::fs::PermissionsExt::from_mode(0o600))
1413        .await
1414        .map_err(|error| format!("secure Codex app-server file '{}': {error}", path.display()))?;
1415    Ok(())
1416}
1417
1418async fn drain_stderr_tail(stderr: tokio::process::ChildStderr, tail: Arc<Mutex<String>>) {
1419    use tokio::io::AsyncBufReadExt;
1420    let mut reader = BufReader::new(stderr);
1421    let mut buffer = Vec::new();
1422    loop {
1423        buffer.clear();
1424        match reader.read_until(b'\n', &mut buffer).await {
1425            Ok(0) | Err(_) => return,
1426            Ok(_) => {
1427                let mut tail = tail.lock().await;
1428                tail.push_str(&String::from_utf8_lossy(&buffer));
1429                if tail.len() > STDERR_TAIL_BYTES {
1430                    let excess = tail.len() - STDERR_TAIL_BYTES;
1431                    let cut = tail
1432                        .char_indices()
1433                        .map(|(index, _)| index)
1434                        .find(|index| *index >= excess)
1435                        .unwrap_or(tail.len());
1436                    tail.drain(..cut);
1437                }
1438            }
1439        }
1440    }
1441}
1442
1443#[cfg(test)]
1444mod tests {
1445    use super::*;
1446    use crate::codex_cli_executor::{
1447        resolve_codex_app_server_permission_config, resolve_codex_auth_config,
1448    };
1449    use bamboo_subagent::executor::HostBridge;
1450    use bamboo_subagent::proto::{PermissionPolicyContext, RunSecrets, SecretValue};
1451
1452    #[cfg(unix)]
1453    fn write_stub_codex(path: &Path) {
1454        use std::os::unix::fs::PermissionsExt as _;
1455        std::fs::write(
1456            path,
1457            r###"#!/bin/sh
1458if [ "$1" = "--version" ]; then
1459  echo 'codex-cli 0.144.5'
1460  exit 0
1461fi
1462if [ "$1" = "exec" ]; then
1463  echo '--json --output-last-message --config --sandbox --dangerously-bypass-approvals-and-sandbox stdin'
1464  exit 0
1465fi
1466if [ "$1" = "app-server" ] && [ "$2" = "--help" ]; then
1467  echo '--listen stdio:// --stdio'
1468  exit 0
1469fi
1470if [ "$1" != "app-server" ]; then
1471  exit 2
1472fi
1473IFS= read -r initialize
1474echo '{"id":1,"result":{"userAgent":"stub/0.144.5"}}'
1475IFS= read -r initialized
1476IFS= read -r thread_start
1477echo '{"id":2,"result":{"thread":{"id":"thread-stub"}}}'
1478IFS= read -r turn_start
1479echo '{"id":3,"result":{"turn":{"id":"turn-stub","status":"inProgress","items":[]}}}'
1480echo '{"id":41,"method":"item/commandExecution/requestApproval","params":{"threadId":"thread-stub","turnId":"turn-stub","itemId":"item-1","command":"touch marker","cwd":"/tmp","reason":"stub command","startedAtMs":1}}'
1481IFS= read -r approval
1482case "$approval" in
1483  *'"decision":"accept"'*) text='stub approved' ;;
1484  *) text='stub denied' ;;
1485esac
1486printf '{"method":"item/agentMessage/delta","params":{"threadId":"thread-stub","turnId":"turn-stub","itemId":"item-2","delta":"%s"}}\n' "$text"
1487echo '{"method":"turn/completed","params":{"threadId":"thread-stub","turn":{"id":"turn-stub","status":"completed","items":[]}}}'
1488while IFS= read -r ignored; do :; done
1489"###,
1490        )
1491        .unwrap();
1492        let mut permissions = std::fs::metadata(path).unwrap().permissions();
1493        permissions.set_mode(0o755);
1494        std::fs::set_permissions(path, permissions).unwrap();
1495    }
1496
1497    #[cfg(unix)]
1498    async fn run_stub(
1499        approved: Option<bool>,
1500        run_auto_approve_permissions: Option<bool>,
1501        provisioned_auto_approve_permissions: bool,
1502    ) -> (ChildOutcome, Vec<Value>) {
1503        let root = tempfile::tempdir().unwrap();
1504        let binary = root.path().join("codex-stub.sh");
1505        write_stub_codex(&binary);
1506        let permissions = resolve_codex_app_server_permission_config(
1507            Some("workspace-write"),
1508            Some("on-request"),
1509            false,
1510            false,
1511            None,
1512            false,
1513            false,
1514        )
1515        .unwrap()
1516        .with_provisioned_permission_resolution(
1517            bamboo_domain::resolve_permission_mode(
1518                if provisioned_auto_approve_permissions {
1519                    bamboo_domain::SessionPermissionMode::Auto
1520                } else {
1521                    bamboo_domain::SessionPermissionMode::Default
1522                },
1523                bamboo_domain::PermissionMode::Default,
1524            ),
1525            "provisioned-stub-session".to_string(),
1526        );
1527        let executor = CodexAppServerExecutor::new(
1528            Some(binary.to_string_lossy().into_owned()),
1529            None,
1530            Some(root.path().to_string_lossy().into_owned()),
1531            Some(root.path().join("state")),
1532            Vec::new(),
1533            CodexAuthConfig::inherit(),
1534            permissions,
1535        )
1536        .await
1537        .unwrap();
1538        let (sink, mut event_rx) = EventSink::channel();
1539        let (sink, approval_task) = if let Some(approved) = approved {
1540            let (host, mut host_rx) = HostBridge::channel();
1541            let task = tokio::spawn(async move {
1542                let request = host_rx.recv().await.expect("approval request");
1543                assert_eq!(request.body["tool_name"], "Bash");
1544                request.reply.send(json!({"approved": approved})).unwrap();
1545            });
1546            (sink.with_host_bridge(host), Some(task))
1547        } else {
1548            (sink, None)
1549        };
1550        let outcome = executor
1551            .run(
1552                RunSpec {
1553                    assignment: "exercise approval".to_string(),
1554                    logical_session: None,
1555                    project_id: None,
1556                    reasoning_effort: None,
1557                    permission_policy: run_auto_approve_permissions.map(
1558                        |auto_approve_permissions| PermissionPolicyContext {
1559                            revision: 1,
1560                            requested_mode: if auto_approve_permissions {
1561                                "auto".to_string()
1562                            } else {
1563                                "default".to_string()
1564                            },
1565                            effective_mode: if auto_approve_permissions {
1566                                "auto".to_string()
1567                            } else {
1568                                "default".to_string()
1569                            },
1570                            bypass_permissions: false,
1571                            auto_approve_permissions,
1572                            session_id: format!("stub-{approved:?}-{auto_approve_permissions}"),
1573                            workspace_path: Some(root.path().to_string_lossy().into_owned()),
1574                            inherit_session_grants: false,
1575                            policy: serde_json::to_value(
1576                                bamboo_tools::permission::SerializablePermissionConfig::default(),
1577                            )
1578                            .unwrap(),
1579                        },
1580                    ),
1581                    messages: Vec::new(),
1582                    activation_run_id: None,
1583                    initial_session_messages: Vec::new(),
1584                    secrets: RunSecrets::default(),
1585                },
1586                sink,
1587                SteerInbox::disconnected(),
1588                CancellationToken::new(),
1589            )
1590            .await;
1591        if let Some(task) = approval_task {
1592            task.await.unwrap();
1593        }
1594        let mut events = Vec::new();
1595        while let Ok(event) = event_rx.try_recv() {
1596            events.push(event);
1597        }
1598        (outcome, events)
1599    }
1600
1601    #[test]
1602    fn recorded_fixture_covers_handshake_and_approval_round_trip() {
1603        let rows =
1604            include_str!("../tests/fixtures/codex-app-server/0.144.5-handshake-approval.jsonl")
1605                .lines()
1606                .map(|line| serde_json::from_str::<Value>(line).expect("valid JSONL row"))
1607                .collect::<Vec<_>>();
1608        assert_eq!(rows[0]["message"]["method"], "initialize");
1609        assert!(rows.iter().any(|row| {
1610            row["direction"] == "client" && row["message"]["method"] == "initialized"
1611        }));
1612        assert!(rows.iter().any(|row| {
1613            row["direction"] == "client" && row["message"]["method"] == "thread/start"
1614        }));
1615        let approval = rows
1616            .iter()
1617            .find(|row| row["message"]["method"] == "item/fileChange/requestApproval")
1618            .expect("approval request");
1619        let approval_id = approval["message"]["id"].clone();
1620        assert!(rows.iter().any(|row| {
1621            row["direction"] == "client"
1622                && row["message"]["id"] == approval_id
1623                && row["message"]["result"]["decision"] == "accept"
1624        }));
1625        assert_eq!(rows.last().unwrap()["message"]["method"], "turn/completed");
1626    }
1627
1628    #[tokio::test]
1629    async fn current_command_approval_relays_allow_to_accept() {
1630        let (write_tx, mut write_rx) = mpsc::unbounded_channel();
1631        let (host, mut host_rx) = HostBridge::channel();
1632        let task = spawn_approval_relay(
1633            write_tx,
1634            json!({
1635                "id": 9,
1636                "method": "item/commandExecution/requestApproval",
1637                "params": {"command": "touch marker", "cwd": "/tmp", "reason": "write marker"}
1638            }),
1639            Some(host),
1640            Duration::from_secs(1),
1641        );
1642        let request = host_rx.recv().await.expect("host approval request");
1643        assert_eq!(request.body["tool_name"], "Bash");
1644        assert_eq!(request.body["resource"], "touch marker");
1645        request
1646            .reply
1647            .send(json!({"approved": true}))
1648            .expect("host reply accepted");
1649        task.await.unwrap();
1650        let response = write_rx.recv().await.expect("app-server response");
1651        assert_eq!(response, json!({"id": 9, "result": {"decision": "accept"}}));
1652    }
1653
1654    #[tokio::test]
1655    async fn approval_timeout_fails_closed_to_decline() {
1656        let (write_tx, mut write_rx) = mpsc::unbounded_channel();
1657        let (host, mut host_rx) = HostBridge::channel();
1658        let task = spawn_approval_relay(
1659            write_tx,
1660            json!({
1661                "id": "approval-10",
1662                "method": "item/fileChange/requestApproval",
1663                "params": {"grantRoot": "/tmp/workspace"}
1664            }),
1665            Some(host),
1666            Duration::from_millis(10),
1667        );
1668        let held_request = host_rx.recv().await.expect("host approval request");
1669        task.await.unwrap();
1670        drop(held_request);
1671        let response = write_rx.recv().await.expect("app-server response");
1672        assert_eq!(
1673            response,
1674            json!({"id": "approval-10", "result": {"decision": "decline"}})
1675        );
1676    }
1677
1678    #[tokio::test]
1679    async fn legacy_apply_patch_denial_uses_legacy_decision_shape() {
1680        let (write_tx, mut write_rx) = mpsc::unbounded_channel();
1681        let task = spawn_approval_relay(
1682            write_tx,
1683            json!({"id": 11, "method": "applyPatchApproval", "params": {}}),
1684            None,
1685            Duration::from_secs(1),
1686        );
1687        task.await.unwrap();
1688        assert_eq!(
1689            write_rx.recv().await.unwrap(),
1690            json!({"id": 11, "result": {"decision": "denied"}})
1691        );
1692    }
1693
1694    #[test]
1695    fn approval_from_another_loaded_thread_is_denied_without_relay() {
1696        let request = json!({
1697            "id": 12,
1698            "method": "item/commandExecution/requestApproval",
1699            "params": {"threadId": "other", "turnId": "turn-stub"}
1700        });
1701        assert!(!approval_matches_active_turn(
1702            &request,
1703            "thread-stub",
1704            "turn-stub"
1705        ));
1706        assert_eq!(
1707            approval_response(&request, false),
1708            json!({"id": 12, "result": {"decision": "decline"}})
1709        );
1710    }
1711
1712    #[cfg(unix)]
1713    #[tokio::test]
1714    async fn subprocess_stub_completes_full_handshake_and_allow_path() {
1715        let (outcome, events) = run_stub(Some(true), Some(false), false).await;
1716        assert_eq!(outcome.result.as_deref(), Some("stub approved"));
1717        assert!(events.iter().any(|event| event["type"] == "complete"));
1718    }
1719
1720    #[cfg(unix)]
1721    #[tokio::test]
1722    async fn subprocess_stub_returns_denial_to_model_and_completes() {
1723        let (outcome, events) = run_stub(Some(false), Some(false), false).await;
1724        assert_eq!(outcome.result.as_deref(), Some("stub denied"));
1725        assert!(events
1726            .iter()
1727            .any(|event| event["type"] == "token" && event["content"] == "stub denied"));
1728    }
1729
1730    #[cfg(unix)]
1731    #[tokio::test]
1732    async fn auto_uses_never_policy_and_denies_unexpected_approval_without_host() {
1733        let (outcome, events) = run_stub(None, Some(true), false).await;
1734        assert_eq!(outcome.result.as_deref(), Some("stub denied"));
1735        assert_eq!(
1736            events.first().and_then(|event| event["type"].as_str()),
1737            Some("permission_posture_activated"),
1738            "typed permission posture must precede every warning/progress event"
1739        );
1740        assert!(events.iter().any(|event| {
1741            event["executor"] == "codex_app_server"
1742                && event["approval_policy"] == "never"
1743                && event["approvals_reviewer"].is_null()
1744                && event["requested_mode"] == "auto"
1745                && event["effective_mode"] == "auto"
1746                && event["executor_mapping"] == "codex_app_server:approvalPolicy=never"
1747        }));
1748    }
1749
1750    #[test]
1751    fn granular_filesystem_and_network_escalation_is_denied_as_empty_subset() {
1752        let request = json!({
1753            "id": 77,
1754            "method": "item/permissions/requestApproval",
1755            "params": {
1756                "threadId": "thread-stub",
1757                "turnId": "turn-stub",
1758                "additionalPermissions": {
1759                    "fileSystem": {"write": ["/outside-workspace"]},
1760                    "network": {"enabled": true}
1761                }
1762            }
1763        });
1764        assert!(is_approval_method("item/permissions/requestApproval"));
1765        assert_eq!(
1766            approval_response(&request, true),
1767            json!({"id": 77, "result": {"permissions": {}}})
1768        );
1769        assert_eq!(
1770            approval_response(&request, false),
1771            json!({"id": 77, "result": {"permissions": {}}})
1772        );
1773    }
1774
1775    #[cfg(unix)]
1776    #[tokio::test]
1777    async fn provisioned_auto_supplies_policy_and_session_when_run_context_is_absent() {
1778        let (outcome, events) = run_stub(None, None, true).await;
1779        assert_eq!(outcome.result.as_deref(), Some("stub denied"));
1780        assert!(events.iter().any(|event| {
1781            event["session_id"] == "provisioned-stub-session"
1782                && event["approval_policy"] == "never"
1783                && event["requested_mode"] == "auto"
1784                && event["executor_mapping"] == "codex_app_server:approvalPolicy=never"
1785        }));
1786    }
1787
1788    #[cfg(unix)]
1789    #[tokio::test]
1790    async fn explicit_run_context_replaces_provisioned_auto_in_app_server() {
1791        let (outcome, events) = run_stub(Some(true), Some(false), true).await;
1792        assert_eq!(outcome.result.as_deref(), Some("stub approved"));
1793        assert!(events.iter().any(|event| {
1794            event["approval_policy"] == "on-request"
1795                && event["approvals_reviewer"] == "user"
1796                && event["requested_mode"] == "default"
1797                && event["executor_mapping"] == "codex_app_server:approvalPolicy=on-request"
1798        }));
1799    }
1800
1801    #[cfg(unix)]
1802    #[tokio::test]
1803    async fn explicit_deny_fails_before_app_server_connection_without_leaking_rule_resource() {
1804        use std::os::unix::fs::PermissionsExt as _;
1805
1806        let root = tempfile::tempdir().unwrap();
1807        let binary = root.path().join("codex-explicit-deny-stub.sh");
1808        let marker = root.path().join("connected");
1809        std::fs::write(
1810            &binary,
1811            r###"#!/bin/sh
1812if [ "$1" = "--version" ]; then
1813  echo 'codex-cli 0.144.5'
1814  exit 0
1815fi
1816if [ "$1" = "exec" ]; then
1817  echo '--json --output-last-message --config --sandbox --dangerously-bypass-approvals-and-sandbox stdin'
1818  exit 0
1819fi
1820if [ "$1" = "app-server" ] && [ "$2" = "--help" ]; then
1821  echo '--listen stdio:// --stdio'
1822  exit 0
1823fi
1824DIR=$(cd "$(dirname "$0")" && pwd)
1825: > "$DIR/connected"
1826exit 2
1827"###,
1828        )
1829        .unwrap();
1830        let mut binary_permissions = std::fs::metadata(&binary).unwrap().permissions();
1831        binary_permissions.set_mode(0o755);
1832        std::fs::set_permissions(&binary, binary_permissions).unwrap();
1833
1834        let permissions = resolve_codex_app_server_permission_config(
1835            Some("workspace-write"),
1836            Some("on-request"),
1837            false,
1838            false,
1839            None,
1840            false,
1841            false,
1842        )
1843        .unwrap();
1844        let executor = CodexAppServerExecutor::new(
1845            Some(binary.to_string_lossy().into_owned()),
1846            None,
1847            Some(root.path().to_string_lossy().into_owned()),
1848            Some(root.path().join("state")),
1849            Vec::new(),
1850            CodexAuthConfig::inherit(),
1851            permissions,
1852        )
1853        .await
1854        .unwrap();
1855        let secret_resource = "TOP_SECRET_CODEX_APP_SERVER_DENY_RESOURCE";
1856        let mut policy = bamboo_tools::permission::SerializablePermissionConfig::default();
1857        policy
1858            .whitelist
1859            .push(bamboo_tools::permission::PermissionRule::new(
1860                bamboo_tools::permission::PermissionType::ExecuteCommand,
1861                secret_resource,
1862                false,
1863            ));
1864        let (sink, mut rx) = EventSink::channel();
1865
1866        let outcome = executor
1867            .run(
1868                RunSpec {
1869                    assignment: "must fail closed".to_string(),
1870                    logical_session: None,
1871                    project_id: None,
1872                    reasoning_effort: None,
1873                    permission_policy: Some(PermissionPolicyContext {
1874                        revision: 29,
1875                        requested_mode: "auto".to_string(),
1876                        effective_mode: "auto".to_string(),
1877                        bypass_permissions: false,
1878                        auto_approve_permissions: true,
1879                        session_id: "codex-app-server-explicit-deny".to_string(),
1880                        workspace_path: Some(root.path().to_string_lossy().into_owned()),
1881                        inherit_session_grants: false,
1882                        policy: serde_json::to_value(policy).unwrap(),
1883                    }),
1884                    messages: Vec::new(),
1885                    activation_run_id: None,
1886                    initial_session_messages: Vec::new(),
1887                    secrets: RunSecrets::default(),
1888                },
1889                sink,
1890                SteerInbox::disconnected(),
1891                CancellationToken::new(),
1892            )
1893            .await;
1894
1895        assert_eq!(
1896            outcome.status,
1897            bamboo_subagent::proto::TerminalStatus::Error
1898        );
1899        assert!(outcome
1900            .error
1901            .as_deref()
1902            .is_some_and(|error| error.contains("explicit-deny")));
1903        assert!(!marker.exists(), "app-server connection must not start");
1904        let events = std::iter::from_fn(|| rx.try_recv().ok()).collect::<Vec<_>>();
1905        assert!(events.iter().any(|event| {
1906            event["type"] == "permission_posture_activated"
1907                && event["executor_mapping"] == "codex_app_server:blocked_explicit_deny"
1908        }));
1909        assert!(events.iter().any(|event| event["type"] == "error"));
1910        assert!(
1911            !serde_json::to_string(&events)
1912                .unwrap()
1913                .contains(secret_resource),
1914            "deny resources must not be emitted in audit/error events"
1915        );
1916    }
1917
1918    #[cfg(unix)]
1919    #[tokio::test]
1920    async fn bamboo_auth_token_file_is_per_run_secret_and_cleared() {
1921        let root = tempfile::tempdir().unwrap();
1922        let binary = root.path().join("codex-stub.sh");
1923        write_stub_codex(&binary);
1924        let auth = resolve_codex_auth_config(
1925            Some("bamboo"),
1926            false,
1927            Some("http://127.0.0.1:9562/openai/v1".to_string()),
1928            Some("responses".to_string()),
1929            None,
1930            &[],
1931            &[],
1932        )
1933        .unwrap();
1934        let permissions = resolve_codex_app_server_permission_config(
1935            Some("workspace-write"),
1936            Some("on-request"),
1937            false,
1938            false,
1939            None,
1940            false,
1941            false,
1942        )
1943        .unwrap();
1944        let executor = CodexAppServerExecutor::new(
1945            Some(binary.to_string_lossy().into_owned()),
1946            None,
1947            Some(root.path().to_string_lossy().into_owned()),
1948            Some(root.path().join("state")),
1949            Vec::new(),
1950            auth,
1951            permissions,
1952        )
1953        .await
1954        .unwrap();
1955        let spec = RunSpec {
1956            assignment: "token lifecycle".to_string(),
1957            logical_session: None,
1958            project_id: None,
1959            reasoning_effort: None,
1960            permission_policy: None,
1961            messages: Vec::new(),
1962            activation_run_id: None,
1963            initial_session_messages: Vec::new(),
1964            secrets: RunSecrets {
1965                codex_provider_token: Some(SecretValue::new("bcx1_app_server_secret")),
1966            },
1967        };
1968        let token_guard = executor
1969            .install_run_token(&spec)
1970            .unwrap()
1971            .expect("bamboo auth installs a token guard");
1972        assert_eq!(
1973            tokio::fs::read_to_string(executor.token_path())
1974                .await
1975                .unwrap(),
1976            "bcx1_app_server_secret"
1977        );
1978        let config = tokio::fs::read_to_string(
1979            executor
1980                .codex_home()
1981                .expect("isolated home")
1982                .join("config.toml"),
1983        )
1984        .await
1985        .unwrap();
1986        assert!(config.contains("codex-provider-token"));
1987        assert!(!config.contains("bcx1_app_server_secret"));
1988        drop(token_guard);
1989        assert!(tokio::fs::read(executor.token_path())
1990            .await
1991            .unwrap()
1992            .is_empty());
1993    }
1994
1995    #[test]
1996    fn logical_session_store_is_bounded_and_preserves_current_session() {
1997        let mut store = AppServerSessionStore::default();
1998        let base = Utc::now();
1999        for index in 0..=MAX_LOGICAL_SESSIONS {
2000            store.sessions.insert(
2001                format!("session-{index}"),
2002                AppServerSessionState {
2003                    thread_id: format!("thread-{index}"),
2004                    workspace: Some("/workspace".to_string()),
2005                    codex_home_mode: "inherit".to_string(),
2006                    updated_at: base + chrono::Duration::seconds(index as i64),
2007                },
2008            );
2009        }
2010
2011        prune_session_store(&mut store, "session-0");
2012
2013        assert_eq!(store.sessions.len(), MAX_LOGICAL_SESSIONS);
2014        assert!(store.sessions.contains_key("session-0"));
2015        assert!(!store.sessions.contains_key("session-1"));
2016    }
2017}