1use 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 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
164pub 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 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 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 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 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 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}