use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use bamboo_agent_core::{AgentError, AgentEvent, Role, Session};
use bamboo_domain::poison::PoisonRecover;
use bamboo_domain::SessionInboxClaim;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use bamboo_subagent::fleet::{spawn_worker_on_bus, SpawnedChild};
use bamboo_subagent::proto::{
AgentRecord, ChildFrame, LogicalSessionIdentity, ParentFrame, PermissionPolicyContext, RunSpec,
SessionMessageDelivery, TerminalStatus,
};
use bamboo_subagent::provision::{
ChildIdentity, ExecutorSpec, ModelRefSpec, Placement, ProvisionSpec, ScopedCredential,
};
use bamboo_subagent::transport::{client_config_trusting_cert, ChildClient};
use crate::runtime::execution::{ExternalChildRunner, SessionInboxRuntimeBinding, SpawnJob};
pub const DEFAULT_MAX_CONCURRENT_ACTORS: usize = 8;
pub const MAX_SPAWN_DEPTH: u32 = 4;
const DEFAULT_MAX_IDLE_PER_KEY: usize = 4;
const POOLED_IDLE_TIMEOUT_SECS: u64 = 300;
const WORKER_FIRST_FRAME_TIMEOUT: Duration = Duration::from_secs(60);
fn active_scoped_session_deny_count(
config: &bamboo_tools::permission::PermissionConfig,
session_id: &str,
) -> usize {
config
.temporary_grants()
.into_iter()
.filter(|grant| {
grant.scope == bamboo_tools::permission::TemporaryPermissionGrantScope::Session
&& grant.effect == bamboo_tools::permission::TemporaryPermissionGrantEffect::Deny
&& grant.session_id.as_deref() == Some(session_id)
})
.count()
}
fn ensure_no_active_scoped_session_denies(
config: &bamboo_tools::permission::PermissionConfig,
session_id: &str,
) -> Result<(), AgentError> {
let count = active_scoped_session_deny_count(config, session_id);
if count == 0 {
Ok(())
} else {
Err(AgentError::LLM(format!(
"external executor activation blocked by {count} active session-scoped explicit deny rule(s)"
)))
}
}
pub struct IssuedCodexRunToken {
pub token_id: String,
pub token: String,
}
pub trait CodexRunTokenAuthority: Send + Sync + 'static {
fn issue(&self, session_id: &str) -> Result<IssuedCodexRunToken, String>;
fn revoke(&self, token_id: &str);
}
struct CodexRunTokenGuard {
authority: Arc<dyn CodexRunTokenAuthority>,
token_id: String,
}
impl Drop for CodexRunTokenGuard {
fn drop(&mut self) {
self.authority.revoke(&self.token_id);
}
}
fn executor_uses_bamboo_codex(executor: &ExecutorSpec) -> bool {
matches!(
executor,
ExecutorSpec::Codex {
auth_mode: Some(mode),
..
} if mode == "bamboo"
) || matches!(
executor,
ExecutorSpec::Codex {
auth_mode: None,
inherit_user_config,
..
} if !inherit_user_config.unwrap_or(false)
)
}
fn executor_has_read_only_permission_profile(executor: &ExecutorSpec) -> bool {
match executor {
ExecutorSpec::ClaudeCode {
permission_mode, ..
} => permission_mode
.as_deref()
.is_some_and(|mode| mode.eq_ignore_ascii_case("plan")),
ExecutorSpec::Codex {
permission_profile,
sandbox,
..
} => {
permission_profile
.as_deref()
.is_some_and(|profile| profile.eq_ignore_ascii_case("read-only"))
|| sandbox
.as_deref()
.is_some_and(|value| value.eq_ignore_ascii_case("read-only"))
}
_ => false,
}
}
fn expected_permission_executor_mapping(
executor: &ExecutorSpec,
resolution: bamboo_domain::PermissionModeResolution,
has_explicit_deny: bool,
) -> Result<Option<String>, AgentError> {
let mapping = match executor {
ExecutorSpec::Echo | ExecutorSpec::CliAdapter { .. } => return Ok(None),
ExecutorSpec::BambooRuntime => {
format!("bamboo_runtime:{}", resolution.effective.as_str())
}
ExecutorSpec::ClaudeCode {
permission_mode, ..
} => {
if has_explicit_deny {
"claude_code:blocked_explicit_deny".to_string()
} else {
let mode = match resolution.effective {
bamboo_domain::PermissionMode::Plan => "plan",
bamboo_domain::PermissionMode::Auto => "bypassPermissions",
bamboo_domain::PermissionMode::AcceptEdits => "acceptEdits",
bamboo_domain::PermissionMode::DontAsk => "dontAsk",
bamboo_domain::PermissionMode::Default
| bamboo_domain::PermissionMode::BypassPermissions => {
permission_mode.as_deref().unwrap_or("default")
}
};
format!("claude_code:permission_mode={mode}")
}
}
ExecutorSpec::Codex {
mode,
sandbox,
approval_policy,
allow_danger_bypass,
..
} => match mode.as_deref().unwrap_or("exec") {
"exec" => {
let approval_policy = expected_codex_exec_approval_policy(
sandbox.as_deref(),
approval_policy.as_deref(),
allow_danger_bypass.unwrap_or(false),
resolution,
)?;
if has_explicit_deny {
"codex_exec:blocked_explicit_deny".to_string()
} else {
format!("codex_exec:approval_policy={approval_policy}")
}
}
"app_server" => {
if !matches!(approval_policy.as_deref(), None | Some("on-request")) {
return Err(AgentError::LLM(
"invalid Codex app-server permission posture configuration".to_string(),
));
}
if has_explicit_deny {
"codex_app_server:blocked_explicit_deny".to_string()
} else {
let approval_policy = if resolution.suppress_approval_prompts()
|| resolution.effective == bamboo_domain::PermissionMode::Plan
{
"never"
} else {
"on-request"
};
format!("codex_app_server:approvalPolicy={approval_policy}")
}
}
_ => {
return Err(AgentError::LLM(
"unsupported Codex executor mode for permission posture contract".to_string(),
));
}
},
};
Ok(Some(mapping))
}
fn expected_codex_exec_approval_policy(
sandbox: Option<&str>,
approval_policy: Option<&str>,
allow_danger_bypass: bool,
resolution: bamboo_domain::PermissionModeResolution,
) -> Result<&'static str, AgentError> {
let configured = match approval_policy {
None | Some("never") => "never",
Some("on-failure") => "on-failure",
Some(_) => {
return Err(AgentError::LLM(
"invalid Codex exec permission posture configuration".to_string(),
));
}
};
if resolution.suppress_approval_prompts()
|| resolution.effective == bamboo_domain::PermissionMode::Plan
{
return Ok("never");
}
match sandbox {
Some("danger-full-access") => Ok("never"),
Some("read-only") | Some("workspace-write") => Ok(configured),
None if allow_danger_bypass || resolution.bypass_permissions() => Ok("never"),
None => Ok(configured),
Some(_) => Err(AgentError::LLM(
"invalid Codex exec permission posture configuration".to_string(),
)),
}
}
fn workspace_is_bamboo_owned(raw: &str) -> bool {
let workspace = std::fs::canonicalize(raw).unwrap_or_else(|_| PathBuf::from(raw));
let configured_root = bamboo_config::paths::resolve_workspace_root();
let configured_root = std::fs::canonicalize(&configured_root).unwrap_or(configured_root);
if workspace.starts_with(&configured_root) {
return true;
}
workspace.ancestors().any(|candidate| {
let Some(name) = candidate.file_name().and_then(|name| name.to_str()) else {
return false;
};
let Some(worktree_root) = candidate.parent() else {
return false;
};
if worktree_root.file_name() != Some(std::ffi::OsStr::new("worktree"))
|| worktree_root.parent().and_then(Path::file_name)
!= Some(std::ffi::OsStr::new(".bamboo"))
{
return false;
}
let marker = worktree_root.join(".bamboo-owned").join(name);
std::fs::read_to_string(marker).is_ok_and(|branch| branch == format!("bamboo/{name}"))
})
}
fn build_codex_run_secrets(
executor: &ExecutorSpec,
authority: Option<Arc<dyn CodexRunTokenAuthority>>,
child_session_id: &str,
) -> Result<
(
bamboo_subagent::proto::RunSecrets,
Option<CodexRunTokenGuard>,
),
AgentError,
> {
if !executor_uses_bamboo_codex(executor) {
return Ok((bamboo_subagent::proto::RunSecrets::default(), None));
}
let authority = authority.ok_or_else(|| {
AgentError::LLM(
"Codex auth mode 'bamboo' requires the server per-run token authority".to_string(),
)
})?;
let issued = authority
.issue(child_session_id)
.map_err(|error| AgentError::LLM(format!("mint Codex per-run provider token: {error}")))?;
let guard = CodexRunTokenGuard {
authority,
token_id: issued.token_id,
};
Ok((
bamboo_subagent::proto::RunSecrets {
codex_provider_token: Some(bamboo_subagent::proto::SecretValue::new(issued.token)),
},
Some(guard),
))
}
struct PooledWorker {
worker: SpawnedChild,
mailbox_id: String,
}
#[derive(Debug, Clone)]
pub struct ResolvedRemotePlacement {
pub endpoint: String,
pub token: Option<String>,
pub ca_cert_file: Option<PathBuf>,
pub host_label: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ResolvedSchedulablePlacement {
pub pool: String,
pub host_label: Option<String>,
}
enum PlacementKind {
Local,
Remote,
Schedulable,
}
pub struct ActorChildRunner {
approval_registry: Option<super::approval_registry::SharedApprovalRegistry>,
permission_config: Option<Arc<bamboo_tools::permission::PermissionConfig>>,
agent_id: String,
worker_bin: PathBuf,
worker_args: Vec<String>,
fabric_dir: PathBuf,
executor: ExecutorSpec,
credentials: Vec<ScopedCredential>,
default_provider: String,
bus: Option<bamboo_subagent::BusEndpoint>,
concurrency: std::sync::Arc<tokio::sync::Semaphore>,
pool: Arc<tokio::sync::Mutex<HashMap<String, Vec<PooledWorker>>>>,
max_idle_per_key: usize,
approval_decider: Option<Arc<dyn ChildApprovalDecider>>,
approval_reviewer: Option<Arc<dyn ChildApprovalReviewer>>,
escalation_bridge: Arc<std::sync::Mutex<Option<bamboo_subagent::executor::HostBridge>>>,
remote_placements: HashMap<String, ResolvedRemotePlacement>,
schedulable_placements: HashMap<String, ResolvedSchedulablePlacement>,
schedule_cursor: Arc<std::sync::Mutex<HashMap<String, usize>>>,
codex_run_tokens: Option<Arc<dyn CodexRunTokenAuthority>>,
session_inbox_runtime: Arc<std::sync::Mutex<Option<SessionInboxRuntimeBinding>>>,
}
#[async_trait]
pub trait ChildApprovalDecider: Send + Sync {
async fn decide(&self, child_session_id: &str, request: &serde_json::Value) -> bool;
}
async fn decide_child_approval(
decider: Option<&Arc<dyn ChildApprovalDecider>>,
child_session_id: &str,
request: &serde_json::Value,
) -> bool {
match decider {
Some(decider) => decider.decide(child_session_id, request).await,
None => false,
}
}
const CHILD_APPROVAL_TIMEOUT: Duration = Duration::from_secs(300);
#[async_trait]
pub trait ChildApprovalReviewer: Send + Sync {
async fn review(
&self,
parent_session_id: &str,
child_session_id: &str,
request: &serde_json::Value,
) -> bool;
}
fn child_approval_reviewer_slot() -> &'static std::sync::OnceLock<Arc<dyn ChildApprovalReviewer>> {
static SLOT: std::sync::OnceLock<Arc<dyn ChildApprovalReviewer>> = std::sync::OnceLock::new();
&SLOT
}
pub fn set_child_approval_reviewer(reviewer: Arc<dyn ChildApprovalReviewer>) {
let _ = child_approval_reviewer_slot().set(reviewer);
}
pub fn child_approval_reviewer() -> Option<Arc<dyn ChildApprovalReviewer>> {
child_approval_reviewer_slot().get().cloned()
}
impl ActorChildRunner {
#[allow(clippy::too_many_arguments)]
pub fn new(
agent_id: String,
worker_bin: PathBuf,
worker_args: Vec<String>,
fabric_dir: PathBuf,
executor: ExecutorSpec,
credentials: Vec<ScopedCredential>,
default_provider: String,
max_concurrent: usize,
) -> Self {
Self {
approval_registry: None,
permission_config: None,
agent_id,
worker_bin,
worker_args,
fabric_dir,
executor,
credentials,
default_provider,
bus: None,
concurrency: std::sync::Arc::new(tokio::sync::Semaphore::new(max_concurrent.max(1))),
pool: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
max_idle_per_key: DEFAULT_MAX_IDLE_PER_KEY,
approval_decider: None,
approval_reviewer: None,
escalation_bridge: Arc::new(std::sync::Mutex::new(None)),
remote_placements: HashMap::new(),
schedulable_placements: HashMap::new(),
schedule_cursor: Arc::new(std::sync::Mutex::new(HashMap::new())),
codex_run_tokens: None,
session_inbox_runtime: Arc::new(std::sync::Mutex::new(None)),
}
}
pub fn with_approval_registry(
mut self,
registry: super::approval_registry::SharedApprovalRegistry,
) -> Self {
self.approval_registry = Some(registry);
self
}
pub fn with_permission_config(
mut self,
config: Arc<bamboo_tools::permission::PermissionConfig>,
) -> Self {
self.permission_config = Some(config);
self
}
pub fn with_bus(mut self, bus: Option<bamboo_subagent::BusEndpoint>) -> Self {
self.bus = bus.filter(|b| !b.endpoint.trim().is_empty());
self
}
pub fn with_approval_decider(mut self, decider: Arc<dyn ChildApprovalDecider>) -> Self {
self.approval_decider = Some(decider);
self
}
pub fn with_approval_reviewer(mut self, reviewer: Arc<dyn ChildApprovalReviewer>) -> Self {
self.approval_reviewer = Some(reviewer);
self
}
pub fn with_codex_run_tokens(
mut self,
authority: Option<Arc<dyn CodexRunTokenAuthority>>,
) -> Self {
self.codex_run_tokens = authority;
self
}
pub fn with_remote_placements(
mut self,
placements: HashMap<String, ResolvedRemotePlacement>,
) -> Self {
self.remote_placements = placements;
self
}
pub fn with_schedulable_placements(
mut self,
placements: HashMap<String, ResolvedSchedulablePlacement>,
) -> Self {
self.schedulable_placements = placements;
self
}
fn fingerprint(spec: &ProvisionSpec) -> String {
let role = spec.identity.role.as_str();
let (provider, model) = spec
.model
.as_ref()
.map(|m| (m.provider.as_str(), m.model.as_str()))
.unwrap_or(("", ""));
let workspace = spec.workspace.as_deref().unwrap_or("");
let mut tools = spec.disabled_tools.clone().unwrap_or_default();
tools.sort();
let caps = &spec.capabilities;
let executor = serde_json::to_string(&spec.executor).unwrap_or_default();
format!(
"{role}\u{1}{provider}\u{1}{model}\u{1}{workspace}\u{1}{}\u{1}d={}\u{1}ns={}\u{1}pr={}\u{1}pe={}\u{1}by={}\u{1}auto={}\u{1}ep={}\u{1}md={}\u{1}nha={}\u{1}gro={}\u{1}executor={executor}",
tools.join(","),
spec.identity.depth,
caps.nested_spawn,
caps.permission_requested_mode,
caps.permission_effective_mode,
caps.bypass,
caps.auto_approve_permissions,
caps.enforce_permissions,
caps.max_spawn_depth.unwrap_or(0),
caps.no_human_approver,
caps.guardian_read_only,
)
}
async fn acquire_bus_worker(
&self,
key: &str,
spec: &ProvisionSpec,
) -> crate::runtime::runner::Result<PooledWorker> {
loop {
let candidate = {
let mut pool = self.pool.lock().await;
pool.get_mut(key).and_then(|bucket| bucket.pop())
};
let Some(mut candidate) = candidate else {
break;
};
if candidate.worker.is_alive() {
return Ok(candidate);
}
candidate.worker.kill().await;
}
let spawned = spawn_worker_on_bus(&self.worker_bin, &self.worker_args, spec)
.await
.map_err(|e| AgentError::LLM(format!("actor spawn (bus) failed: {e}")))?;
let mailbox_id = spawned.record.agent_id.clone();
Ok(PooledWorker {
worker: spawned,
mailbox_id,
})
}
async fn release_bus_worker(&self, key: &str, mut worker: PooledWorker) {
if !worker.worker.is_alive() {
worker.worker.kill().await;
return;
}
let mut pool = self.pool.lock().await;
let bucket = pool.entry(key.to_string()).or_default();
if bucket.len() >= self.max_idle_per_key {
drop(pool);
worker.worker.kill().await;
return;
}
bucket.push(worker);
}
fn build_spec(&self, session: &Session, job: &SpawnJob) -> ProvisionSpec {
let mut spec = ProvisionSpec::new(
ChildIdentity {
child_id: job.child_session_id.clone(),
parent_id: Some(job.parent_session_id.clone()),
project_key: None,
role: session
.metadata
.get("subagent_type")
.cloned()
.unwrap_or_else(|| "worker".to_string()),
depth: session.spawn_depth,
},
self.executor.clone(),
self.fabric_dir.to_string_lossy().into_owned(),
);
spec.workspace = session.workspace.clone();
if let ExecutorSpec::Codex {
workspace_owned, ..
} = &mut spec.executor
{
*workspace_owned = Some(
spec.workspace
.as_deref()
.is_some_and(workspace_is_bamboo_owned),
);
}
spec.bus = self.bus.clone();
spec.model = session
.model_ref
.as_ref()
.map(|r| ModelRefSpec {
provider: r.provider.clone(),
model: r.model.clone(),
})
.or_else(|| {
let m = job.model.trim();
(!m.is_empty()).then(|| ModelRefSpec {
provider: self.default_provider.clone(),
model: m.to_string(),
})
});
spec.disabled_tools = job.disabled_tools.clone();
match &spec.executor {
ExecutorSpec::Codex {
auth_mode,
provider_key_ref,
..
} => {
if auth_mode.as_deref() == Some("custom") {
if let Some(reference) = provider_key_ref {
if let Some(credential) = self.credentials.iter().find(|credential| {
credential.credential_ref.as_deref() == Some(reference)
}) {
spec.secrets.provider_credentials.push(credential.clone());
} else {
tracing::warn!(
"actor child {}: custom Codex credential reference '{}' did not resolve",
job.child_session_id,
reference
);
}
} else {
tracing::warn!(
"actor child {}: custom Codex executor has no credential reference",
job.child_session_id
);
}
}
}
_ => {
let provider = spec
.model
.as_ref()
.map(|model| model.provider.as_str())
.filter(|provider| !provider.trim().is_empty())
.unwrap_or(&self.default_provider);
if let Some(credential) = self
.credentials
.iter()
.find(|credential| credential.provider == provider)
{
spec.secrets.provider_credentials.push(credential.clone());
} else {
tracing::warn!(
"actor child {}: no credential found for provider '{}'",
job.child_session_id,
provider
);
}
}
}
spec.capabilities.nested_spawn = session.spawn_depth < MAX_SPAWN_DEPTH;
spec.capabilities.max_spawn_depth = Some(MAX_SPAWN_DEPTH);
spec.capabilities.enforce_permissions = true;
let requested_permission_mode = session
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
.unwrap_or_default();
let configured_permission_mode = self
.permission_config
.as_ref()
.map(|config| config.mode())
.unwrap_or_default();
spec.capabilities.no_human_approver = session
.agent_runtime_state
.as_ref()
.is_some_and(|s| s.no_human_approver);
spec.capabilities.guardian_read_only =
session.metadata.get("subagent_type").map(String::as_str) == Some("guardian");
if spec.capabilities.guardian_read_only {
if let ExecutorSpec::Codex {
permission_profile, ..
} = &mut spec.executor
{
*permission_profile = Some("read-only".to_string());
}
}
let read_only_overlay = session
.agent_runtime_state
.as_ref()
.is_some_and(|state| state.plan_mode.is_some())
|| spec.capabilities.guardian_read_only
|| executor_has_read_only_permission_profile(&spec.executor);
let permission_resolution = bamboo_domain::resolve_permission_mode_with_read_only(
requested_permission_mode,
configured_permission_mode,
read_only_overlay,
);
spec.capabilities.bypass = permission_resolution.bypass_permissions();
spec.capabilities.auto_approve_permissions =
permission_resolution.suppress_approval_prompts();
spec.capabilities.permission_requested_mode =
permission_resolution.requested.as_str().to_string();
spec.capabilities.permission_effective_mode =
permission_resolution.effective.as_str().to_string();
if let Some(placement) = self.remote_placements.get(spec.identity.role.as_str()) {
spec.placement = Placement::Remote {
endpoint: placement.endpoint.clone(),
};
spec.secrets.worker_auth_token = placement.token.clone();
} else if let Some(placement) = self.schedulable_placements.get(spec.identity.role.as_str())
{
spec.placement = Placement::Schedulable {
pool: placement.pool.clone(),
};
}
spec
}
fn placement_stamp_for(&self, spec: &ProvisionSpec) -> Option<String> {
let host_label = match &spec.placement {
Placement::Remote { .. } => self
.remote_placements
.get(spec.identity.role.as_str())
.and_then(|p| p.host_label.as_deref()),
Placement::Schedulable { .. } => self
.schedulable_placements
.get(spec.identity.role.as_str())
.and_then(|p| p.host_label.as_deref()),
Placement::Local => None,
};
placement_metadata(&spec.placement, host_label)
}
async fn resolve_schedulable_worker(
&self,
role: &str,
) -> std::result::Result<String, AgentError> {
let pool = self
.schedulable_placements
.get(role)
.ok_or_else(|| {
AgentError::LLM(format!(
"schedulable placement for role '{role}' vanished before scheduling"
))
})?
.pool
.clone();
let bus = self.bus.as_ref().ok_or_else(|| {
AgentError::LLM(format!(
"schedulable role '{role}': no mailbox bus configured (subagents.broker)"
))
})?;
let mut q = bamboo_broker::BrokerClient::connect(
&bus.endpoint,
bamboo_subagent::AgentRef {
session_id: format!("sched-q-{role}"),
role: None,
},
&bus.token,
)
.await
.map_err(|e| {
AgentError::LLM(format!(
"schedulable role '{role}': bus connect failed: {e}"
))
})?;
let candidates = q.list_connected(&pool).await.map_err(|e| {
AgentError::LLM(format!(
"schedulable role '{role}': bus presence query failed: {e}"
))
})?;
if candidates.is_empty() {
return Err(AgentError::LLM(format!(
"schedulable role '{role}': no live worker in pool '{pool}' on the bus \
(NOT spawning a local subprocess — a schedulable role has no local fallback)"
)));
}
let idx = {
let mut cursors = self.schedule_cursor.lock().recover_poison();
let cursor = cursors.entry(pool.clone()).or_insert(0);
let i = *cursor % candidates.len();
*cursor = cursor.wrapping_add(1);
i
};
Ok(candidates[idx].clone())
}
}
#[async_trait]
impl ExternalChildRunner for ActorChildRunner {
async fn should_handle(&self, session: &Session) -> bool {
session.metadata.get("runtime.kind") == Some(&"external".to_string())
&& session.metadata.get("external.protocol") == Some(&"actor".to_string())
&& session.metadata.get("external.agent_id") == Some(&self.agent_id)
}
fn set_escalation_bridge(&self, bridge: Option<bamboo_subagent::executor::HostBridge>) {
*self.escalation_bridge.lock().recover_poison() = bridge;
}
fn set_session_inbox_runtime(&self, binding: Option<SessionInboxRuntimeBinding>) {
*self.session_inbox_runtime.lock().recover_poison() = binding;
}
async fn execute_external_child(
&self,
session: &mut Session,
job: &SpawnJob,
event_tx: mpsc::Sender<AgentEvent>,
cancel_token: CancellationToken,
) -> crate::runtime::runner::Result<()> {
let escalation = self.escalation_bridge.lock().recover_poison().clone();
let session_inbox_runtime = self.session_inbox_runtime.lock().recover_poison().clone();
let assignment = extract_assignment(session);
let mut spec = self.build_spec(session, job);
spec.reusable = true;
if spec.limits.idle_timeout_secs.is_none() {
spec.limits.idle_timeout_secs = Some(POOLED_IDLE_TIMEOUT_SECS);
}
let pool_key = Self::fingerprint(&spec);
if executor_uses_bamboo_codex(&spec.executor) && !matches!(spec.placement, Placement::Local)
{
return Err(AgentError::LLM(
"Codex auth mode 'bamboo' requires local actor placement; use custom mode with a reachable URL for remote workers"
.to_string(),
));
}
let project_id = project_id_for_actor_run(session)?;
let requested_permission_mode = session
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
.unwrap_or_default();
let permission_policy = if let Some(config) = self.permission_config.as_ref() {
ensure_no_active_scoped_session_denies(config, &session.id)?;
let read_only_overlay = session
.agent_runtime_state
.as_ref()
.is_some_and(|state| state.plan_mode.is_some())
|| spec.capabilities.guardian_read_only
|| executor_has_read_only_permission_profile(&spec.executor);
let resolution = bamboo_domain::resolve_permission_mode_with_read_only(
requested_permission_mode,
config.mode(),
read_only_overlay,
);
let policy = serde_json::to_value(config.to_serializable()).map_err(|error| {
AgentError::LLM(format!(
"failed to serialize permission policy for external executor: {error}"
))
})?;
Some(PermissionPolicyContext {
revision: config.policy_revision(),
requested_mode: resolution.requested.as_str().to_string(),
effective_mode: resolution.effective.as_str().to_string(),
bypass_permissions: resolution.bypass_permissions(),
auto_approve_permissions: resolution.suppress_approval_prompts(),
session_id: session.id.clone(),
workspace_path: session.workspace.clone(),
inherit_session_grants: false,
policy,
})
} else {
None
};
let provisioned_permission =
spec.capabilities.permission_resolution().map_err(|error| {
AgentError::LLM(format!("invalid provisioned permission posture: {error}"))
})?;
let policy_resolution = permission_policy
.as_ref()
.map(|context| {
context.resolved_modes().map(|(requested, effective)| {
bamboo_domain::PermissionModeResolution {
requested,
effective,
}
})
})
.transpose()
.map_err(|error| {
AgentError::LLM(format!("invalid host permission posture: {error}"))
})?;
let host_audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata);
let has_explicit_deny = self.permission_config.as_ref().is_some_and(|config| {
bamboo_tools::permission::explicit_deny_policy_reason(&config.to_serializable())
.is_some()
});
let expected_executor_mapping = expected_permission_executor_mapping(
&spec.executor,
policy_resolution.unwrap_or(provisioned_permission),
has_explicit_deny,
)?;
let expected_permission_posture =
expected_executor_mapping.map(|executor_mapping| ExpectedPermissionPosture {
policy_revision: permission_policy
.as_ref()
.map(|context| context.revision)
.or_else(|| host_audit.as_ref().map(|audit| audit.policy_revision))
.unwrap_or_default(),
resolution: policy_resolution.unwrap_or(provisioned_permission),
expected_audit_revision: host_audit.as_ref().map(|audit| audit.audit_revision),
executor_mapping,
});
let _slot = self
.concurrency
.acquire()
.await
.map_err(|_| AgentError::LLM("actor concurrency limiter closed".to_string()))?;
let (run_secrets, _codex_token_guard) = build_codex_run_secrets(
&spec.executor,
self.codex_run_tokens.clone(),
&job.child_session_id,
)?;
let kind = match spec.placement {
Placement::Remote { .. } => PlacementKind::Remote,
Placement::Schedulable { .. } => PlacementKind::Schedulable,
Placement::Local => PlacementKind::Local,
};
let remote = !matches!(kind, PlacementKind::Local);
if let Some(placement_meta) = self.placement_stamp_for(&spec) {
session
.metadata
.insert("placement".to_string(), placement_meta);
}
let mut attempt = 0u8;
let (result, actor) = loop {
let (actor, mut client) = match kind {
PlacementKind::Remote => {
let placement = self
.remote_placements
.get(spec.identity.role.as_str())
.ok_or_else(|| {
AgentError::LLM(format!(
"remote placement for role '{}' vanished before connect",
spec.identity.role
))
})?;
let endpoint = placement.endpoint.clone();
let trust_cfg = match placement.ca_cert_file.as_deref() {
Some(path) => Some(client_config_trusting_cert(path).map_err(|e| {
AgentError::LLM(format!(
"remote worker CA cert '{}': {e}",
path.display()
))
})?),
None => None,
};
let client = ChildClient::connect_with_auth_tls(
&endpoint,
placement.token.as_deref(),
trust_cfg,
)
.await
.map_err(|e| {
AgentError::LLM(format!("remote actor connect to '{endpoint}' failed: {e}"))
})?;
let record = AgentRecord {
agent_id: job.child_session_id.clone(),
role: spec.identity.role.clone(),
labels: Vec::new(),
endpoint: endpoint.clone(),
pid: 0,
version: String::new(),
started_at: chrono::Utc::now(),
lease_expires_at: chrono::Utc::now(),
};
let _ = endpoint;
let actor = PooledWorker {
worker: SpawnedChild::remote(record),
mailbox_id: job.child_session_id.clone(),
};
let client: Box<dyn bamboo_subagent::ChildLink> = Box::new(client);
(actor, client)
}
PlacementKind::Schedulable => {
let bus = self.bus.as_ref().ok_or_else(|| {
AgentError::LLM(
"schedulable sub-agents require a mailbox bus (subagents.broker)"
.to_string(),
)
})?;
let mailbox_id = self
.resolve_schedulable_worker(spec.identity.role.as_str())
.await?;
let parent = bamboo_subagent::AgentRef {
session_id: format!("p-{}", job.child_session_id),
role: None,
};
let link = bamboo_broker::BrokerChildLink::connect(
&bus.endpoint,
parent,
&bus.token,
mailbox_id.clone(),
)
.await
.map_err(|e| {
AgentError::LLM(format!(
"schedulable link connect to '{mailbox_id}' failed: {e}"
))
})?;
let actor = PooledWorker {
worker: SpawnedChild::remote(AgentRecord {
agent_id: mailbox_id.clone(),
role: spec.identity.role.clone(),
labels: Vec::new(),
endpoint: bus.endpoint.clone(),
pid: 0,
version: String::new(),
started_at: chrono::Utc::now(),
lease_expires_at: chrono::Utc::now(),
}),
mailbox_id,
};
let client: Box<dyn bamboo_subagent::ChildLink> = Box::new(link);
(actor, client)
}
PlacementKind::Local => {
let bus = self.bus.as_ref().ok_or_else(|| {
AgentError::LLM(
"local sub-agents require a mailbox bus (subagents.broker); none is \
configured and the bus could not be embedded"
.to_string(),
)
})?;
let actor = self.acquire_bus_worker(&pool_key, &spec).await?;
let parent = bamboo_subagent::AgentRef {
session_id: format!("p-{}", job.child_session_id),
role: None,
};
let link = bamboo_broker::BrokerChildLink::connect(
&bus.endpoint,
parent,
&bus.token,
actor.mailbox_id.clone(),
)
.await
.map_err(|e| {
AgentError::LLM(format!("broker child link connect failed: {e}"))
})?;
let client: Box<dyn bamboo_subagent::ChildLink> = Box::new(link);
(actor, client)
}
};
let (delivery_tx, mut delivery_rx) = mpsc::unbounded_channel::<u64>();
let bound_activation_run_id = match session_inbox_runtime.as_ref() {
Some(binding) => {
let run_id = binding
.router
.attach_delivery_sink(&job.child_session_id, delivery_tx.clone())
.await;
if run_id.is_none() {
tracing::debug!(
session_id = %job.child_session_id,
"actor driver had no current SessionInbox activation owner to bind"
);
}
run_id
}
None => None,
};
drop(delivery_tx);
let initial_pairs = match (
session_inbox_runtime.as_ref(),
bound_activation_run_id.as_deref(),
) {
(Some(binding), Some(run_id)) => {
match claim_canonical_deliveries(binding, session, run_id, usize::MAX).await {
Ok(deliveries) => deliveries,
Err(error) => {
binding
.router
.detach_delivery_sink(&job.child_session_id, run_id)
.await;
if !remote {
actor.worker.kill().await;
}
return Err(error);
}
}
}
_ => Vec::new(),
};
let initial_session_messages = initial_pairs
.iter()
.map(|(_, delivery)| delivery.clone())
.collect::<Vec<_>>();
let initial_inflight_claims = initial_pairs
.into_iter()
.map(|(claim, _)| claim)
.collect::<VecDeque<_>>();
let messages = session
.messages
.iter()
.filter_map(|message| serde_json::to_value(message).ok())
.collect();
if let Err(e) = client
.send(ParentFrame::Run(RunSpec {
assignment: assignment.clone(),
logical_session: Some(logical_identity_for_actor_run(session, job)),
project_id: project_id.clone(),
reasoning_effort: None,
permission_policy: permission_policy.clone(),
messages,
activation_run_id: bound_activation_run_id.clone(),
initial_session_messages,
secrets: run_secrets.clone(),
}))
.await
{
if let (Some(binding), Some(run_id)) = (
session_inbox_runtime.as_ref(),
bound_activation_run_id.as_deref(),
) {
binding
.router
.detach_delivery_sink(&job.child_session_id, run_id)
.await;
}
if !remote {
actor.worker.kill().await;
}
return Err(AgentError::LLM(format!("actor run dispatch failed: {e}")));
}
let (live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
let live_guard = super::live::register(
&job.child_session_id,
live_tx,
attempt as u32,
self.approval_registry.clone(),
);
let result = drive(ActorDriveContext {
client: &mut *client,
parent_session_id: &job.parent_session_id,
child_session_id: &job.child_session_id,
child_attempt: attempt as u32,
approval_registry: self.approval_registry.as_ref(),
approval_decider: self.approval_decider.as_ref(),
approval_reviewer: self.approval_reviewer.as_ref(),
escalation_bridge: escalation.clone(),
event_tx: &event_tx,
cancel_token: &cancel_token,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: session,
expected_permission_posture: expected_permission_posture.clone(),
session_inbox_runtime: session_inbox_runtime.as_ref(),
activation_run_id: bound_activation_run_id.as_deref(),
initial_inflight_claims,
first_frame_timeout: Some(WORKER_FIRST_FRAME_TIMEOUT),
})
.await;
if let (Some(binding), Some(run_id)) = (
session_inbox_runtime.as_ref(),
bound_activation_run_id.as_deref(),
) {
binding
.router
.detach_delivery_sink(&job.child_session_id, run_id)
.await;
}
drop(live_guard);
drop(client);
if attempt == 0 && matches!(result, Err(AgentError::WorkerUnresponsive(_))) {
match kind {
PlacementKind::Local => {
tracing::warn!(
"actor child {} got no first frame; reaping the worker and respawning once",
job.child_session_id
);
actor.worker.kill().await;
attempt += 1;
continue;
}
PlacementKind::Schedulable => {
tracing::warn!(
"scheduled actor child {} got no first frame; re-selecting a pool worker",
job.child_session_id
);
drop(actor);
attempt += 1;
continue;
}
PlacementKind::Remote => {}
}
}
break (result, actor);
};
if remote {
drop(actor);
} else {
match &result {
Ok(_) => self.release_bus_worker(&pool_key, actor).await,
Err(_) => actor.worker.kill().await,
}
}
match result {
Ok(Some(text)) => {
if !text.is_empty() {
session.add_message(bamboo_agent_core::Message::assistant(text, None));
}
Ok(())
}
Ok(None) => Ok(()),
Err(e) => Err(e),
}
}
}
fn placement_metadata(placement: &Placement, host_label: Option<&str>) -> Option<String> {
let value = match placement {
Placement::Local => return None,
Placement::Remote { endpoint } => serde_json::json!({
"kind": "remote",
"host": host_label.map(str::to_string).unwrap_or_else(|| host_of_endpoint(endpoint)),
}),
Placement::Schedulable { pool } => serde_json::json!({
"kind": "remote",
"host": host_label.unwrap_or(pool),
}),
};
serde_json::to_string(&value).ok()
}
fn host_of_endpoint(endpoint: &str) -> String {
endpoint
.trim()
.trim_start_matches("wss://")
.trim_start_matches("ws://")
.split(['/', ':'])
.next()
.unwrap_or(endpoint)
.to_string()
}
async fn reconcile_already_admitted_claim(
binding: &SessionInboxRuntimeBinding,
session: &mut Session,
claim: &SessionInboxClaim,
) -> crate::runtime::runner::Result<()> {
let latest = binding
.storage
.load_session(&session.id)
.await
.map_err(|error| {
AgentError::LLM(format!(
"load canonical SessionInbox checkpoint for {}: {error}",
session.id
))
})?
.ok_or_else(|| {
AgentError::LLM(format!(
"canonical SessionInbox target disappeared: {}",
session.id
))
})?;
let Some(message) = latest
.messages
.iter()
.find(|message| bamboo_domain::is_matching_session_message(message, &claim.envelope))
.cloned()
else {
return Err(AgentError::LLM(format!(
"canonical admitted receipt for {} exists without transcript message {}",
session.id, claim.envelope.id
)));
};
if let Some(existing) = session
.messages
.iter_mut()
.find(|existing| existing.id == message.id)
{
*existing = message;
} else {
session.add_message(message);
}
bamboo_domain::merge_session_inbox_admission(session, &latest);
binding
.inbox
.ack(&session.id, claim)
.await
.map_err(|error| {
AgentError::LLM(format!(
"ack recovered canonical SessionInbox claim {}: {error}",
claim.envelope.id
))
})
}
async fn checkpoint_and_ack_canonical_claim(
binding: &SessionInboxRuntimeBinding,
session: &mut Session,
claim: &SessionInboxClaim,
) -> crate::runtime::runner::Result<()> {
if claim.envelope.target_session_id != session.id {
return Err(AgentError::LLM(format!(
"canonical SessionInbox claim target {} does not match active logical session {}",
claim.envelope.target_session_id, session.id
)));
}
if binding
.inbox
.was_admitted(&session.id, &claim.envelope.id)
.await
.map_err(|error| {
AgentError::LLM(format!(
"inspect canonical admitted receipt {}: {error}",
claim.envelope.id
))
})?
{
return reconcile_already_admitted_claim(binding, session, claim).await;
}
let transcript_has_id = session
.messages
.iter()
.any(|message| bamboo_domain::is_matching_session_message(message, &claim.envelope));
if session
.messages
.iter()
.any(|message| message.id == claim.envelope.id.as_str())
&& !transcript_has_id
{
return Err(AgentError::LLM(format!(
"canonical SessionInbox id {} collides with a non-matching transcript message",
claim.envelope.id
)));
}
let cursor_has_id = session
.session_inbox_admission()
.is_some_and(|state| state.contains(&claim.envelope.id));
if cursor_has_id && !transcript_has_id {
return Err(AgentError::LLM(format!(
"canonical SessionInbox cursor exists without transcript message {}",
claim.envelope.id
)));
}
let before = session.clone();
if !transcript_has_id {
let message = claim.envelope.to_provider_message().map_err(|error| {
AgentError::LLM(format!(
"translate canonical SessionInbox envelope {}: {error}",
claim.envelope.id
))
})?;
session.add_message(message);
}
session
.session_inbox_admission_mut()
.record(claim.envelope.id.clone(), claim.generation);
session.updated_at = chrono::Utc::now();
if let Err(error) = binding
.persistence
.checkpoint_runtime_session(session)
.await
{
*session = before;
return Err(AgentError::LLM(format!(
"checkpoint canonical SessionInbox claim {}: {error}",
claim.envelope.id
)));
}
if !session
.messages
.iter()
.any(|message| bamboo_domain::is_matching_session_message(message, &claim.envelope))
{
*session = before;
return Err(AgentError::LLM(format!(
"canonical SessionInbox checkpoint lost typed transcript proof for {}",
claim.envelope.id
)));
}
binding
.inbox
.ack(&session.id, claim)
.await
.map_err(|error| {
AgentError::LLM(format!(
"ack canonical SessionInbox claim {} after checkpoint: {error}",
claim.envelope.id
))
})
}
async fn checkpoint_claim_context_before_dispatch(
binding: &SessionInboxRuntimeBinding,
session: &mut Session,
claims: &[SessionInboxClaim],
) -> crate::runtime::runner::Result<()> {
if claims.is_empty() {
return Ok(());
}
let before = session.clone();
for claim in claims {
if claim.envelope.target_session_id != session.id {
return Err(AgentError::LLM(format!(
"canonical SessionInbox claim target {} does not match actor session {}",
claim.envelope.target_session_id, session.id
)));
}
let matching = session
.messages
.iter()
.any(|message| bamboo_domain::is_matching_session_message(message, &claim.envelope));
if session
.messages
.iter()
.any(|message| message.id == claim.envelope.id.as_str())
&& !matching
{
return Err(AgentError::LLM(format!(
"canonical SessionInbox id {} collides before actor dispatch",
claim.envelope.id
)));
}
if session
.session_inbox_admission()
.is_some_and(|state| state.contains(&claim.envelope.id))
&& !matching
{
return Err(AgentError::LLM(format!(
"canonical SessionInbox cursor exists without transcript proof for {}",
claim.envelope.id
)));
}
if !matching {
let message = claim.envelope.to_provider_message().map_err(|error| {
AgentError::LLM(format!(
"translate canonical SessionInbox envelope {} before actor dispatch: {error}",
claim.envelope.id
))
})?;
session.add_message(message);
}
}
session.updated_at = chrono::Utc::now();
if let Err(error) = binding
.persistence
.checkpoint_runtime_session(session)
.await
{
*session = before;
return Err(AgentError::LLM(format!(
"checkpoint canonical SessionInbox actor context: {error}"
)));
}
for claim in claims {
if !session
.messages
.iter()
.any(|message| bamboo_domain::is_matching_session_message(message, &claim.envelope))
{
*session = before;
return Err(AgentError::LLM(format!(
"actor context checkpoint lost typed transcript proof for {}",
claim.envelope.id
)));
}
}
Ok(())
}
async fn claim_canonical_deliveries(
binding: &SessionInboxRuntimeBinding,
session: &mut Session,
activation_run_id: &str,
limit: usize,
) -> crate::runtime::runner::Result<Vec<(SessionInboxClaim, SessionMessageDelivery)>> {
let claims = binding
.inbox
.claim(&session.id, limit)
.await
.map_err(|error| {
AgentError::LLM(format!(
"claim canonical SessionInbox for active actor {}: {error}",
session.id
))
})?;
if claims.is_empty() {
return Ok(Vec::new());
}
let interrupt_generation = binding
.inbox
.inspect(&session.id)
.await
.map_err(|error| {
AgentError::LLM(format!(
"inspect canonical SessionInbox activation policy for {}: {error}",
session.id
))
})?
.interrupt_generation;
let mut unconfirmed = Vec::with_capacity(claims.len());
for claim in claims {
if binding
.inbox
.was_admitted(&session.id, &claim.envelope.id)
.await
.map_err(|error| {
AgentError::LLM(format!(
"inspect canonical SessionInbox claim {}: {error}",
claim.envelope.id
))
})?
{
reconcile_already_admitted_claim(binding, session, &claim).await?;
continue;
}
unconfirmed.push(claim);
}
checkpoint_claim_context_before_dispatch(binding, session, &unconfirmed).await?;
let mut deliveries = Vec::with_capacity(unconfirmed.len());
for claim in unconfirmed {
if session
.session_inbox_admission()
.is_some_and(|state| state.contains(&claim.envelope.id))
{
binding
.inbox
.ack(&session.id, &claim)
.await
.map_err(|error| {
AgentError::LLM(format!(
"finish confirmed canonical SessionInbox ack {}: {error}",
claim.envelope.id
))
})?;
continue;
}
let activation_policy = if claim.generation <= interrupt_generation {
bamboo_domain::SessionActivationPolicy::InterruptSpecificWait
} else {
bamboo_domain::SessionActivationPolicy::RespectSpecificWait
};
let delivery = SessionMessageDelivery {
target_session_id: session.id.clone(),
envelope: claim.envelope.clone(),
canonical_claim_generation: claim.generation,
activation_run_id: activation_run_id.to_string(),
activation_policy,
};
deliveries.push((claim, delivery));
}
Ok(deliveries)
}
async fn forward_next_canonical_claim(
client: &mut dyn bamboo_subagent::ChildLink,
binding: &SessionInboxRuntimeBinding,
session: &mut Session,
activation_run_id: &str,
inflight: &mut VecDeque<SessionInboxClaim>,
) -> crate::runtime::runner::Result<()> {
if !inflight.is_empty() {
return Ok(());
}
let Some((claim, delivery)) =
claim_canonical_deliveries(binding, session, activation_run_id, 1)
.await?
.pop()
else {
return Ok(());
};
client
.send(ParentFrame::SessionMessage { delivery })
.await
.map_err(|error| {
AgentError::LLM(format!(
"forward canonical SessionInbox claim {} to active actor: {error}",
claim.envelope.id
))
})?;
inflight.push_back(claim);
Ok(())
}
struct ActorDriveContext<'a> {
client: &'a mut dyn bamboo_subagent::ChildLink,
parent_session_id: &'a str,
child_session_id: &'a str,
child_attempt: u32,
approval_registry: Option<&'a super::approval_registry::SharedApprovalRegistry>,
approval_decider: Option<&'a Arc<dyn ChildApprovalDecider>>,
approval_reviewer: Option<&'a Arc<dyn ChildApprovalReviewer>>,
escalation_bridge: Option<bamboo_subagent::executor::HostBridge>,
event_tx: &'a mpsc::Sender<AgentEvent>,
cancel_token: &'a CancellationToken,
live_rx: &'a mut mpsc::UnboundedReceiver<ParentFrame>,
delivery_rx: &'a mut mpsc::UnboundedReceiver<u64>,
logical_session: &'a mut Session,
expected_permission_posture: Option<ExpectedPermissionPosture>,
session_inbox_runtime: Option<&'a SessionInboxRuntimeBinding>,
activation_run_id: Option<&'a str>,
initial_inflight_claims: VecDeque<SessionInboxClaim>,
first_frame_timeout: Option<Duration>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ExpectedPermissionPosture {
policy_revision: u64,
resolution: bamboo_domain::PermissionModeResolution,
expected_audit_revision: Option<u64>,
executor_mapping: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PermissionPostureHandshake {
NotRequired,
Awaiting,
Confirmed,
}
impl PermissionPostureHandshake {
fn new(expected: Option<&ExpectedPermissionPosture>) -> Self {
if expected.is_some() {
Self::Awaiting
} else {
Self::NotRequired
}
}
fn is_awaiting(self) -> bool {
self == Self::Awaiting
}
fn posture_was_confirmed(self) -> bool {
self == Self::Confirmed
}
}
fn permission_posture_seed_from_event(
session: &Session,
event: &AgentEvent,
) -> Result<Option<bamboo_domain::PermissionAuditSeed>, String> {
let AgentEvent::PermissionPostureActivated {
session_id,
policy_revision,
requested_mode,
effective_mode,
executor_mapping,
} = event
else {
return Ok(None);
};
if session_id != &session.id {
return Err("permission posture event targets a different logical session".to_string());
}
let requested = bamboo_domain::SessionPermissionMode::from_audit_str(requested_mode)
.ok_or_else(|| "permission posture event has an invalid requested mode".to_string())?;
let effective = bamboo_domain::PermissionMode::from_audit_str(effective_mode)
.ok_or_else(|| "permission posture event has an invalid effective mode".to_string())?;
let resolution = bamboo_domain::PermissionModeResolution {
requested,
effective,
};
if !resolution.is_consistent() {
return Err("permission posture event has an inconsistent mode pair".to_string());
}
let current_requested = session
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
.unwrap_or_default();
if current_requested != requested {
return Err("permission posture event is stale for the host typed mode".to_string());
}
let mapping_chars = executor_mapping.chars().count();
if mapping_chars == 0 || mapping_chars > bamboo_domain::MAX_PERMISSION_EXECUTOR_MAPPING_CHARS {
return Err("permission posture event has an invalid executor mapping".to_string());
}
Ok(Some(bamboo_domain::PermissionAuditSeed::new(
*policy_revision,
resolution,
executor_mapping,
)))
}
async fn drive(context: ActorDriveContext<'_>) -> crate::runtime::runner::Result<Option<String>> {
let ActorDriveContext {
client,
parent_session_id,
child_session_id,
child_attempt,
approval_registry,
approval_decider,
approval_reviewer,
escalation_bridge,
event_tx,
cancel_token,
live_rx,
delivery_rx,
logical_session,
expected_permission_posture,
session_inbox_runtime,
activation_run_id,
initial_inflight_claims,
first_frame_timeout,
} = context;
let mut got_first_frame = false;
let mut first_frame_watch = first_frame_timeout.map(|d| Box::pin(tokio::time::sleep(d)));
let mut inflight_claims = initial_inflight_claims;
let strict_permission_events = expected_permission_posture.is_some();
let mut permission_handshake =
PermissionPostureHandshake::new(expected_permission_posture.as_ref());
loop {
tokio::select! {
_ = cancel_token.cancelled() => {
break;
}
_ = async {
match first_frame_watch.as_mut() {
Some(s) => s.as_mut().await,
None => std::future::pending::<()>().await,
}
}, if !got_first_frame => {
return Err(AgentError::WorkerUnresponsive(format!(
"child {child_session_id} produced no frame within {:?}",
first_frame_timeout.unwrap_or_default()
)));
}
Some(_generation) = delivery_rx.recv(),
if session_inbox_runtime.is_some() && activation_run_id.is_some() =>
{
forward_next_canonical_claim(
client,
session_inbox_runtime.expect("guarded"),
logical_session,
activation_run_id.expect("guarded"),
&mut inflight_claims,
)
.await?;
}
Some(frame) = live_rx.recv() => {
if client.send(frame).await.is_err() {
tracing::warn!("live steering frame could not be sent; connection failing");
}
}
frame = client.next_frame() => {
got_first_frame = true;
first_frame_watch = None;
match frame {
Ok(Some(ChildFrame::Event { event })) => {
let ev = match serde_json::from_value::<AgentEvent>(event) {
Ok(ev) => ev,
Err(error) if strict_permission_events => {
return Err(AgentError::LLM(format!(
"actor emitted malformed AgentEvent under a typed permission posture contract: {error}"
)));
}
Err(_) => continue,
};
if matches!(&ev, AgentEvent::PermissionPostureActivated { .. }) {
if permission_handshake.posture_was_confirmed() {
return Err(AgentError::LLM(
"actor emitted a duplicate permission posture activation"
.to_string(),
));
}
let seed = permission_posture_seed_from_event(logical_session, &ev)
.map_err(AgentError::LLM)?
.ok_or_else(|| {
AgentError::LLM(
"actor permission posture event did not decode as a posture"
.to_string(),
)
})?;
if let Some(expected) = expected_permission_posture.as_ref() {
if seed.policy_revision != expected.policy_revision
|| seed.resolution != expected.resolution
{
return Err(AgentError::LLM(
"permission posture event does not match the host-dispatched policy"
.to_string(),
));
}
if seed.executor_mapping() != expected.executor_mapping {
return Err(AgentError::LLM(
"permission posture event does not match the host-dispatched executor mapping"
.to_string(),
));
}
}
if let Some(binding) = session_inbox_runtime {
let saved = binding
.persistence
.record_permission_posture_activation(
&logical_session.id,
expected_permission_posture
.as_ref()
.and_then(|expected| expected.expected_audit_revision),
&seed,
)
.await
.map_err(|error| {
AgentError::LLM(format!(
"persist child permission posture bootstrap: {error}"
))
})?
.ok_or_else(|| {
AgentError::LLM(
"persist child permission posture bootstrap: session not found"
.to_string(),
)
})?;
let snapshot = bamboo_domain::PermissionAuditSnapshot::from_metadata(
&saved.metadata,
)
.ok_or_else(|| {
AgentError::LLM(
"persisted child permission posture audit is incomplete"
.to_string(),
)
})?;
snapshot.write_to(&mut logical_session.metadata);
} else {
bamboo_domain::record_permission_audit(
&mut logical_session.metadata,
&seed,
None,
)
.map_err(|error| {
AgentError::LLM(format!(
"record in-memory child permission posture: {error}"
))
})?;
}
permission_handshake = PermissionPostureHandshake::Confirmed;
} else if permission_handshake.is_awaiting() {
return Err(AgentError::LLM(
"actor emitted an execution event before permission posture confirmation"
.to_string(),
));
}
let _ = event_tx.send(ev).await;
}
Ok(Some(ChildFrame::ApprovalRequest { id, body })) => {
if permission_handshake.is_awaiting() {
return Err(AgentError::LLM(
"actor requested approval before permission posture confirmation"
.to_string(),
));
}
if let Some(reviewer) = approval_reviewer
.cloned()
.or_else(child_approval_reviewer)
{
let child = child_session_id.to_string();
let parent = parent_session_id.to_string();
let req_id = id.clone();
let body = body.clone();
let registry = approval_registry.cloned();
tokio::spawn(async move {
let approved = tokio::time::timeout(
CHILD_APPROVAL_TIMEOUT,
reviewer.review(&parent, &child, &body),
)
.await
.unwrap_or(false);
super::live::deliver_approval_scoped(
registry.as_ref(),
&child,
child_attempt,
&req_id,
approved,
);
});
} else if approval_decider.is_some() {
let approved =
decide_child_approval(approval_decider, child_session_id, &body)
.await;
if client
.send(ParentFrame::ApprovalReply { id, approved })
.await
.is_err()
{
tracing::warn!(
"failed to answer approval_request; connection failing"
);
}
} else if let Some(host) = escalation_bridge.clone() {
let child = child_session_id.to_string();
let req_id = id.clone();
let body = body.clone();
let registry = approval_registry.cloned();
tokio::spawn(async move {
let approved = match tokio::time::timeout(
CHILD_APPROVAL_TIMEOUT,
host.approval_call(body),
)
.await
{
Ok(Ok(reply)) => reply
.get("approved")
.and_then(|v| v.as_bool())
.unwrap_or(false),
_ => false,
};
super::live::deliver_approval_scoped(
registry.as_ref(),
&child,
child_attempt,
&req_id,
approved,
);
});
} else {
tracing::warn!(
parent_session_id,
child_session_id,
request_id = %id,
"forced-ask request has no parent-agent reviewer; denying"
);
if client
.send(ParentFrame::ApprovalReply {
id,
approved: false,
})
.await
.is_err()
{
tracing::warn!(
"failed to send fail-closed approval reply; connection failing"
);
}
}
}
Ok(Some(ChildFrame::SessionMessageAdmitted { confirmation })) => {
let Some(binding) = session_inbox_runtime else {
tracing::warn!(
child_session_id,
"ignoring SessionInbox confirmation without a runtime binding"
);
continue;
};
let Some(bound_run_id) = activation_run_id else {
tracing::warn!(
child_session_id,
"ignoring SessionInbox confirmation without an activation owner"
);
continue;
};
let Some(claim) = inflight_claims.front() else {
tracing::warn!(
child_session_id,
envelope_id = %confirmation.envelope_id,
"rejecting stale SessionInbox confirmation with no in-flight canonical claim"
);
continue;
};
let exact = confirmation.target_session_id == logical_session.id
&& confirmation.envelope_id == claim.envelope.id.as_str()
&& confirmation.canonical_claim_generation == claim.generation
&& confirmation.activation_run_id == bound_run_id;
if !exact
|| !binding
.router
.owns_run(&logical_session.id, bound_run_id)
.await
{
tracing::warn!(
child_session_id,
expected_target = %logical_session.id,
received_target = %confirmation.target_session_id,
expected_envelope_id = %claim.envelope.id,
received_envelope_id = %confirmation.envelope_id,
expected_generation = claim.generation,
received_generation = confirmation.canonical_claim_generation,
expected_run_id = bound_run_id,
received_run_id = %confirmation.activation_run_id,
"rejecting stale or mismatched SessionInbox admission confirmation"
);
continue;
}
let claim = inflight_claims
.pop_front()
.expect("validated in-flight canonical claim");
checkpoint_and_ack_canonical_claim(binding, logical_session, &claim)
.await?;
if inflight_claims.is_empty() {
forward_next_canonical_claim(
client,
binding,
logical_session,
bound_run_id,
&mut inflight_claims,
)
.await?;
}
}
Ok(Some(ChildFrame::Terminal { status, result, error, .. })) => {
if permission_handshake.is_awaiting() {
return Err(AgentError::LLM(
"actor terminated before permission posture confirmation"
.to_string(),
));
}
if let Some(claim) = inflight_claims.front() {
return Err(AgentError::LLM(format!(
"actor terminated before durably admitting SessionInbox message {}; canonical claim remains recoverable",
claim.envelope.id
)));
}
return match status {
TerminalStatus::Completed => Ok(result),
TerminalStatus::Cancelled => Err(AgentError::Cancelled),
TerminalStatus::Error => Err(AgentError::LLM(
error.unwrap_or_else(|| "actor child errored".to_string()),
)),
TerminalStatus::Suspended => Err(AgentError::LLM(
"nested sub-agent suspend received but resume transport is not wired"
.to_string(),
)),
};
}
Ok(None) => {
return Err(AgentError::LLM(
"actor child closed before terminal".to_string(),
));
}
Err(e) => {
return Err(AgentError::LLM(format!("actor transport error: {e}")));
}
}
}
}
}
let _ = client.send(ParentFrame::Cancel).await;
Err(AgentError::Cancelled)
}
fn project_id_for_actor_run(
session: &Session,
) -> Result<Option<bamboo_domain::ProjectId>, AgentError> {
match crate::project_context::ProjectContextResolver::session_project_identity(session) {
crate::project_context::SessionProjectIdentity::Assigned(project_id) => {
Ok(Some(project_id))
}
crate::project_context::SessionProjectIdentity::Unassigned => Ok(None),
crate::project_context::SessionProjectIdentity::Invalid { raw, message } => {
Err(AgentError::LLM(format!(
"child session carries an invalid Project identity '{raw}': {message}"
)))
}
}
}
fn logical_identity_for_actor_run(session: &Session, job: &SpawnJob) -> LogicalSessionIdentity {
LogicalSessionIdentity {
session_id: session.id.clone(),
parent_session_id: session
.parent_session_id
.clone()
.or_else(|| Some(job.parent_session_id.clone())),
root_session_id: if session.root_session_id.trim().is_empty() {
job.parent_session_id.clone()
} else {
session.root_session_id.clone()
},
}
}
fn extract_assignment(session: &Session) -> String {
session
.messages
.iter()
.rev()
.find(|m| matches!(m.role, Role::User))
.map(|m| m.content.clone())
.unwrap_or_else(|| {
session
.metadata
.get("title")
.cloned()
.unwrap_or_else(|| "Execute task".to_string())
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::SessionActivationRouter;
use bamboo_domain::{RuntimeSessionPersistence, SessionInboxPort, Storage};
#[test]
fn actor_preflight_counts_only_current_session_scoped_denies() {
let config = bamboo_tools::permission::PermissionConfig::new();
let secret_matcher = "TOP_SECRET_ACTOR_DENY_MATCHER";
config.deny_scoped_session_permission(
"target-session",
bamboo_tools::permission::PermissionType::ExecuteCommand,
secret_matcher,
);
config.deny_scoped_session_permission(
"other-session",
bamboo_tools::permission::PermissionType::WriteFile,
"/other/**",
);
assert_eq!(
active_scoped_session_deny_count(&config, "target-session"),
1
);
assert_eq!(
active_scoped_session_deny_count(&config, "clean-session"),
0
);
let error = ensure_no_active_scoped_session_denies(&config, "target-session")
.unwrap_err()
.to_string();
assert!(!error.contains(secret_matcher));
ensure_no_active_scoped_session_denies(&config, "clean-session")
.expect("another session's deny must not block this activation");
}
#[test]
fn permission_posture_event_rejects_oversized_executor_mapping() {
let session = Session::new("mapping-bound", "model");
let event = AgentEvent::PermissionPostureActivated {
session_id: session.id.clone(),
policy_revision: 1,
requested_mode: "default".to_string(),
effective_mode: "default".to_string(),
executor_mapping: "x".repeat(bamboo_domain::MAX_PERMISSION_EXECUTOR_MAPPING_CHARS + 1),
};
assert!(permission_posture_seed_from_event(&session, &event)
.unwrap_err()
.contains("executor mapping"));
}
#[test]
fn remote_audit_revision_and_timestamp_fields_cannot_poison_host_audit() {
let mut session = Session::new("host-resigns-audit", "model");
let hostile_timestamp = "9".repeat(1024);
let event: AgentEvent = serde_json::from_value(serde_json::json!({
"type": "permission_posture_activated",
"session_id": session.id.clone(),
"policy_revision": 31,
"requested_mode": "default",
"effective_mode": "default",
"executor_mapping": "codex_exec:approval_policy=never",
"audit_revision": u64::MAX,
"transitioned_at": hostile_timestamp,
}))
.expect("unknown remote audit fields are ignored by the typed event");
let seed = permission_posture_seed_from_event(&session, &event)
.unwrap()
.expect("permission event");
let host_revision =
bamboo_domain::record_permission_audit(&mut session.metadata, &seed, None).unwrap();
let host_audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata)
.expect("host-generated complete audit");
assert_eq!(host_audit.audit_revision, host_revision);
assert!(host_audit.audit_revision < bamboo_domain::MAX_PERMISSION_AUDIT_REVISION);
assert_ne!(host_audit.transitioned_at, hostile_timestamp);
assert!(chrono::DateTime::parse_from_rfc3339(&host_audit.transitioned_at).is_ok());
}
struct ActorFaultingPersistence {
inner: Arc<bamboo_storage::LockedSessionStore>,
fail_checkpoint_once: std::sync::atomic::AtomicBool,
}
#[async_trait]
impl RuntimeSessionPersistence for ActorFaultingPersistence {
async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
self.inner.merge_save_runtime(session).await
}
async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
if self
.fail_checkpoint_once
.swap(false, std::sync::atomic::Ordering::SeqCst)
{
return Err(std::io::Error::other("injected actor checkpoint failure"));
}
self.inner.checkpoint_runtime_session(session).await
}
async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
self.inner.storage().load_session(session_id).await
}
}
struct ActorFailBeforeAckInbox {
inner: Arc<dyn SessionInboxPort>,
fail_once: std::sync::atomic::AtomicBool,
}
#[async_trait]
impl SessionInboxPort for ActorFailBeforeAckInbox {
async fn deliver(
&self,
envelope: &bamboo_domain::SessionMessageEnvelope,
) -> Result<bamboo_domain::SessionInboxReceipt, bamboo_domain::SessionInboxError> {
self.inner.deliver(envelope).await
}
async fn mark_activation_eligible(
&self,
target_session_id: &str,
generation: u64,
policy: bamboo_domain::SessionActivationPolicy,
) -> Result<(), bamboo_domain::SessionInboxError> {
self.inner
.mark_activation_eligible(target_session_id, generation, policy)
.await
}
async fn claim(
&self,
target_session_id: &str,
limit: usize,
) -> Result<Vec<SessionInboxClaim>, bamboo_domain::SessionInboxError> {
self.inner.claim(target_session_id, limit).await
}
async fn was_admitted(
&self,
target_session_id: &str,
id: &bamboo_domain::SessionMessageId,
) -> Result<bool, bamboo_domain::SessionInboxError> {
self.inner.was_admitted(target_session_id, id).await
}
async fn ack(
&self,
target_session_id: &str,
claim: &SessionInboxClaim,
) -> Result<(), bamboo_domain::SessionInboxError> {
if self
.fail_once
.swap(false, std::sync::atomic::Ordering::SeqCst)
{
return Err(bamboo_domain::SessionInboxError::Storage(
"injected actor pre-ack failure".to_string(),
));
}
self.inner.ack(target_session_id, claim).await
}
async fn inspect(
&self,
target_session_id: &str,
) -> Result<bamboo_domain::SessionInboxBacklog, bamboo_domain::SessionInboxError> {
self.inner.inspect(target_session_id).await
}
}
async fn actor_inbox_fixture(
session_id: &str,
) -> (
tempfile::TempDir,
Arc<bamboo_storage::SessionStoreV2>,
Arc<bamboo_storage::LockedSessionStore>,
Arc<dyn SessionInboxPort>,
Session,
SessionInboxClaim,
) {
let temp = tempfile::tempdir().unwrap();
let store = Arc::new(
bamboo_storage::SessionStoreV2::new(temp.path().to_path_buf())
.await
.unwrap(),
);
let storage: Arc<dyn Storage> = store.clone();
let locked = Arc::new(bamboo_storage::LockedSessionStore::new(storage));
let inbox: Arc<dyn SessionInboxPort> = Arc::new(bamboo_storage::FileSessionInbox::new(
store.clone(),
bamboo_domain::SessionInboxLimits::default(),
));
let session = Session::new(session_id, "model");
store.save_session(&session).await.unwrap();
let mut envelope =
bamboo_domain::SessionMessageEnvelope::user_input(session_id, "actor follow-up");
envelope.id =
bamboo_domain::SessionMessageId::parse(format!("{session_id}-message")).unwrap();
let receipt = inbox.deliver(&envelope).await.unwrap();
inbox
.mark_activation_eligible(
session_id,
receipt.generation,
bamboo_domain::SessionActivationPolicy::InterruptSpecificWait,
)
.await
.unwrap();
let claim = inbox.claim(session_id, 1).await.unwrap().remove(0);
(temp, store, locked, inbox, session, claim)
}
fn actor_binding(
store: Arc<bamboo_storage::SessionStoreV2>,
inbox: Arc<dyn SessionInboxPort>,
persistence: Arc<dyn RuntimeSessionPersistence>,
) -> SessionInboxRuntimeBinding {
let storage: Arc<dyn Storage> = store;
SessionInboxRuntimeBinding {
router: SessionActivationRouter::new(),
inbox,
storage,
persistence,
}
}
#[tokio::test]
async fn actor_mismatched_typed_marker_id_collision_never_acks() {
let (_temp, store, locked, inbox, mut session, claim) =
actor_inbox_fixture("actor-live-id-collision").await;
let mut forged = claim.envelope.to_provider_message().unwrap();
forged.metadata = Some(serde_json::json!({
"session_message": {
"id": claim.envelope.id,
"target_session_id": "different-session"
}
}));
session.add_message(forged);
let persistence: Arc<dyn RuntimeSessionPersistence> = locked;
let binding = actor_binding(store, inbox.clone(), persistence);
assert!(
checkpoint_and_ack_canonical_claim(&binding, &mut session, &claim)
.await
.is_err()
);
assert_eq!(inbox.inspect(&session.id).await.unwrap().claimed, 1);
assert!(!inbox
.was_admitted(&session.id, &claim.envelope.id)
.await
.unwrap());
}
#[tokio::test]
async fn actor_concurrent_durable_id_collision_after_claim_never_acks() {
let (_temp, store, locked, inbox, mut session, claim) =
actor_inbox_fixture("actor-durable-id-collision").await;
let mut concurrent = store.load_session(&session.id).await.unwrap().unwrap();
let mut forged = bamboo_agent_core::Message::user("concurrent actor collision");
forged.id = claim.envelope.id.to_string();
concurrent.add_message(forged);
store.save_session(&concurrent).await.unwrap();
let persistence: Arc<dyn RuntimeSessionPersistence> = locked;
let binding = actor_binding(store.clone(), inbox.clone(), persistence);
assert!(
checkpoint_and_ack_canonical_claim(&binding, &mut session, &claim)
.await
.is_err()
);
assert_eq!(inbox.inspect(&session.id).await.unwrap().claimed, 1);
assert!(!inbox
.was_admitted(&session.id, &claim.envelope.id)
.await
.unwrap());
let durable = store.load_session(&session.id).await.unwrap().unwrap();
assert!(!durable.messages.iter().any(|message| {
bamboo_domain::is_matching_session_message(message, &claim.envelope)
}));
}
#[tokio::test]
async fn actor_concurrent_durable_typed_body_mismatch_never_acks() {
let (_temp, store, locked, inbox, mut session, claim) =
actor_inbox_fixture("actor-durable-typed-body-collision").await;
let mut different = claim.envelope.clone();
different.body = bamboo_domain::SessionMessageBody::Content(
bamboo_domain::SessionMessageContent::text("forged actor body"),
);
let mut concurrent = store.load_session(&session.id).await.unwrap().unwrap();
concurrent.add_message(different.to_provider_message().unwrap());
store.save_session(&concurrent).await.unwrap();
let persistence: Arc<dyn RuntimeSessionPersistence> = locked;
let binding = actor_binding(store.clone(), inbox.clone(), persistence);
assert!(
checkpoint_and_ack_canonical_claim(&binding, &mut session, &claim)
.await
.is_err()
);
assert_eq!(inbox.inspect(&session.id).await.unwrap().claimed, 1);
assert!(!inbox
.was_admitted(&session.id, &claim.envelope.id)
.await
.unwrap());
let durable = store.load_session(&session.id).await.unwrap().unwrap();
assert!(!durable
.messages
.iter()
.any(|message| bamboo_domain::is_matching_session_message(message, &claim.envelope)));
}
#[tokio::test]
async fn actor_checkpoint_failure_rolls_back_and_restart_admits_once() {
let (_temp, store, locked, inbox, mut session, claim) =
actor_inbox_fixture("actor-checkpoint-failure").await;
let envelope_id = claim.envelope.id.clone();
let fault: Arc<dyn RuntimeSessionPersistence> = Arc::new(ActorFaultingPersistence {
inner: locked.clone(),
fail_checkpoint_once: std::sync::atomic::AtomicBool::new(true),
});
let binding = actor_binding(store.clone(), inbox.clone(), fault);
assert!(
checkpoint_and_ack_canonical_claim(&binding, &mut session, &claim)
.await
.is_err()
);
assert!(!session
.messages
.iter()
.any(|message| message.id == envelope_id.as_str()));
assert_eq!(inbox.inspect(&session.id).await.unwrap().claimed, 1);
assert!(!inbox.was_admitted(&session.id, &envelope_id).await.unwrap());
let reopened: Arc<dyn SessionInboxPort> = Arc::new(bamboo_storage::FileSessionInbox::new(
store.clone(),
bamboo_domain::SessionInboxLimits::default(),
));
let recovered = reopened.claim(&session.id, 1).await.unwrap().remove(0);
let persistence: Arc<dyn RuntimeSessionPersistence> = locked;
let binding = actor_binding(store.clone(), reopened.clone(), persistence);
let mut restarted = store.load_session(&session.id).await.unwrap().unwrap();
checkpoint_and_ack_canonical_claim(&binding, &mut restarted, &recovered)
.await
.unwrap();
assert_eq!(
restarted
.messages
.iter()
.filter(|message| message.id == envelope_id.as_str())
.count(),
1
);
assert!(reopened
.was_admitted(&session.id, &envelope_id)
.await
.unwrap());
let backlog = reopened.inspect(&session.id).await.unwrap();
assert_eq!(backlog.pending + backlog.claimed, 0);
}
#[tokio::test]
async fn actor_checkpoint_success_pre_ack_failure_recovers_without_duplicate() {
let (_temp, store, locked, real_inbox, mut session, claim) =
actor_inbox_fixture("actor-pre-ack-failure").await;
let envelope_id = claim.envelope.id.clone();
let faulted: Arc<dyn SessionInboxPort> = Arc::new(ActorFailBeforeAckInbox {
inner: real_inbox.clone(),
fail_once: std::sync::atomic::AtomicBool::new(true),
});
let persistence: Arc<dyn RuntimeSessionPersistence> = locked.clone();
let binding = actor_binding(store.clone(), faulted, persistence);
assert!(
checkpoint_and_ack_canonical_claim(&binding, &mut session, &claim)
.await
.is_err()
);
let durable = store.load_session(&session.id).await.unwrap().unwrap();
assert_eq!(
durable
.messages
.iter()
.filter(|message| message.id == envelope_id.as_str())
.count(),
1
);
assert_eq!(real_inbox.inspect(&session.id).await.unwrap().claimed, 1);
assert!(!real_inbox
.was_admitted(&session.id, &envelope_id)
.await
.unwrap());
let reopened: Arc<dyn SessionInboxPort> = Arc::new(bamboo_storage::FileSessionInbox::new(
store.clone(),
bamboo_domain::SessionInboxLimits::default(),
));
let recovered = reopened.claim(&session.id, 1).await.unwrap().remove(0);
let persistence: Arc<dyn RuntimeSessionPersistence> = locked;
let binding = actor_binding(store.clone(), reopened.clone(), persistence);
let mut restarted = durable;
checkpoint_and_ack_canonical_claim(&binding, &mut restarted, &recovered)
.await
.unwrap();
assert_eq!(
restarted
.messages
.iter()
.filter(|message| message.id == envelope_id.as_str())
.count(),
1
);
assert!(reopened
.was_admitted(&session.id, &envelope_id)
.await
.unwrap());
let backlog = reopened.inspect(&session.id).await.unwrap();
assert_eq!(backlog.pending + backlog.claimed, 0);
}
struct ConfirmationSequenceLink {
frames: VecDeque<ChildFrame>,
sent: Vec<ParentFrame>,
}
#[async_trait]
impl bamboo_subagent::ChildLink for ConfirmationSequenceLink {
async fn send(&mut self, frame: ParentFrame) -> bamboo_subagent::TransportResult<()> {
self.sent.push(frame);
Ok(())
}
async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
match self.frames.pop_front() {
Some(frame) => Ok(Some(frame)),
None => std::future::pending().await,
}
}
}
fn admission_confirmation(
session_id: &str,
claim: &SessionInboxClaim,
run_id: &str,
) -> bamboo_subagent::proto::SessionMessageAdmissionConfirmation {
bamboo_subagent::proto::SessionMessageAdmissionConfirmation {
target_session_id: session_id.to_string(),
envelope_id: claim.envelope.id.to_string(),
canonical_claim_generation: claim.generation,
activation_run_id: run_id.to_string(),
}
}
fn expected_default_permission_posture(policy_revision: u64) -> ExpectedPermissionPosture {
ExpectedPermissionPosture {
policy_revision,
resolution: bamboo_domain::PermissionModeResolution {
requested: bamboo_domain::SessionPermissionMode::Default,
effective: bamboo_domain::PermissionMode::Default,
},
expected_audit_revision: None,
executor_mapping: "test_actor:permission_mode=default".to_string(),
}
}
fn permission_posture_frame(session_id: &str, policy_revision: u64) -> ChildFrame {
ChildFrame::Event {
event: serde_json::to_value(AgentEvent::PermissionPostureActivated {
session_id: session_id.to_string(),
policy_revision,
requested_mode: "default".to_string(),
effective_mode: "default".to_string(),
executor_mapping: "test_actor:permission_mode=default".to_string(),
})
.expect("serialize permission posture event"),
}
}
fn actor_event_frame(event: AgentEvent) -> ChildFrame {
ChildFrame::Event {
event: serde_json::to_value(event).expect("serialize actor event"),
}
}
fn completed_actor_frame() -> ChildFrame {
ChildFrame::Terminal {
status: TerminalStatus::Completed,
result: Some("done".to_string()),
error: None,
transcript: Vec::new(),
}
}
async fn drive_permission_handshake_frames(
session_id: &str,
frames: impl IntoIterator<Item = ChildFrame>,
expected: ExpectedPermissionPosture,
) -> (
crate::runtime::runner::Result<Option<String>>,
Session,
Vec<AgentEvent>,
Vec<ParentFrame>,
) {
let mut link = ConfirmationSequenceLink {
frames: frames.into_iter().collect(),
sent: Vec::new(),
};
let (event_tx, mut event_rx) = mpsc::channel(16);
let cancel = CancellationToken::new();
let (_live_tx, mut live_rx) = mpsc::unbounded_channel();
let (_delivery_tx, mut delivery_rx) = mpsc::unbounded_channel();
let mut session = Session::new(session_id, "model");
let result = drive(ActorDriveContext {
client: &mut link,
parent_session_id: "permission-parent",
child_session_id: session_id,
child_attempt: 0,
approval_registry: None,
approval_decider: None,
approval_reviewer: None,
escalation_bridge: None,
event_tx: &event_tx,
cancel_token: &cancel,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: &mut session,
expected_permission_posture: Some(expected),
session_inbox_runtime: None,
activation_run_id: None,
initial_inflight_claims: VecDeque::new(),
first_frame_timeout: Some(Duration::from_secs(1)),
})
.await;
let events = std::iter::from_fn(|| event_rx.try_recv().ok()).collect();
(result, session, events, link.sent)
}
#[tokio::test]
async fn actor_permission_handshake_rejects_terminal_without_posture() {
let session_id = "permission-missing";
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[completed_actor_frame()],
expected_default_permission_posture(7),
)
.await;
assert!(result
.unwrap_err()
.to_string()
.contains("terminated before permission posture confirmation"));
assert!(events.is_empty());
assert!(bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata).is_none());
}
#[tokio::test]
async fn actor_permission_handshake_rejects_malformed_agent_event() {
let session_id = "permission-malformed";
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[ChildFrame::Event {
event: serde_json::json!({"type": "token", "content": 42}),
}],
expected_default_permission_posture(7),
)
.await;
assert!(result
.unwrap_err()
.to_string()
.contains("malformed AgentEvent"));
assert!(events.is_empty());
assert!(bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata).is_none());
}
#[tokio::test]
async fn actor_permission_handshake_rejects_execution_event_before_posture() {
let session_id = "permission-early-event";
let early_events = [
(
"progress",
AgentEvent::RunnerProgress {
session_id: session_id.to_string(),
round_count: 1,
},
),
(
"token",
AgentEvent::Token {
content: "must-not-forward".to_string(),
},
),
(
"tool",
AgentEvent::ToolStart {
tool_call_id: "early-tool".to_string(),
tool_name: "Read".to_string(),
arguments: serde_json::json!({"file_path": "README.md"}),
},
),
];
for (kind, event) in early_events {
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[
actor_event_frame(event),
permission_posture_frame(session_id, 7),
completed_actor_frame(),
],
expected_default_permission_posture(7),
)
.await;
assert!(
result
.unwrap_err()
.to_string()
.contains("execution event before permission posture confirmation"),
"{kind} must fail closed before posture"
);
assert!(events.is_empty(), "{kind} must not be forwarded");
assert!(
bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata).is_none(),
"{kind} must not advance the permission audit"
);
}
}
#[tokio::test]
async fn actor_permission_handshake_rejects_approval_before_posture() {
let session_id = "permission-early-approval";
let (result, session, events, sent) = drive_permission_handshake_frames(
session_id,
[ChildFrame::ApprovalRequest {
id: "approval-before-posture".to_string(),
body: serde_json::json!({
"tool_name": "Bash",
"permission": "execute",
"resource": "echo must-not-run"
}),
}],
expected_default_permission_posture(7),
)
.await;
assert!(result
.unwrap_err()
.to_string()
.contains("requested approval before permission posture confirmation"));
assert!(events.is_empty());
assert!(
sent.is_empty(),
"an unconfirmed actor must receive no approval reply"
);
assert!(bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata).is_none());
}
#[tokio::test]
async fn actor_permission_handshake_rejects_mismatched_posture() {
let session_id = "permission-mismatch";
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[permission_posture_frame(session_id, 8)],
expected_default_permission_posture(7),
)
.await;
assert!(result
.unwrap_err()
.to_string()
.contains("does not match the host-dispatched policy"));
assert!(events.is_empty());
assert!(bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata).is_none());
}
#[tokio::test]
async fn actor_permission_handshake_rejects_untrusted_executor_mapping() {
let session_id = "permission-hostile-mapping";
for hostile_mapping in [
"wrong_executor:permission_mode=default",
"test_actor:permission_mode=default;credential=must-not-persist",
] {
let frame = ChildFrame::Event {
event: serde_json::to_value(AgentEvent::PermissionPostureActivated {
session_id: session_id.to_string(),
policy_revision: 7,
requested_mode: "default".to_string(),
effective_mode: "default".to_string(),
executor_mapping: hostile_mapping.to_string(),
})
.expect("serialize hostile posture fixture"),
};
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[frame],
expected_default_permission_posture(7),
)
.await;
let error = result.unwrap_err().to_string();
assert!(error.contains("host-dispatched executor mapping"));
assert!(
!error.contains(hostile_mapping),
"untrusted mapping must not be reflected in host errors"
);
assert!(events.is_empty());
assert!(
bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata).is_none(),
"untrusted mapping must not reach durable or in-memory audit state"
);
}
}
#[tokio::test]
async fn actor_permission_handshake_rejects_duplicate_posture() {
let session_id = "permission-duplicate";
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[
permission_posture_frame(session_id, 7),
permission_posture_frame(session_id, 7),
completed_actor_frame(),
],
expected_default_permission_posture(7),
)
.await;
assert!(result
.unwrap_err()
.to_string()
.contains("duplicate permission posture activation"));
assert_eq!(
events.len(),
1,
"only the confirmed posture may be forwarded"
);
assert!(matches!(
events[0],
AgentEvent::PermissionPostureActivated { .. }
));
let audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata)
.expect("the first matching posture must be recorded");
assert_eq!(audit.policy_revision, 7);
}
#[tokio::test]
async fn actor_permission_handshake_happy_path_persists_before_forwarding_execution() {
let session_id = "permission-happy";
let (result, session, events, _) = drive_permission_handshake_frames(
session_id,
[
permission_posture_frame(session_id, 7),
actor_event_frame(AgentEvent::RunnerProgress {
session_id: session_id.to_string(),
round_count: 1,
}),
actor_event_frame(AgentEvent::Token {
content: "working".to_string(),
}),
actor_event_frame(AgentEvent::ToolStart {
tool_call_id: "tool-1".to_string(),
tool_name: "Read".to_string(),
arguments: serde_json::json!({"file_path": "README.md"}),
}),
completed_actor_frame(),
],
expected_default_permission_posture(7),
)
.await;
assert_eq!(result.unwrap().as_deref(), Some("done"));
assert!(matches!(
events.as_slice(),
[
AgentEvent::PermissionPostureActivated { .. },
AgentEvent::RunnerProgress { .. },
AgentEvent::Token { .. },
AgentEvent::ToolStart { .. }
]
));
let audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&session.metadata)
.expect("matching posture must be recorded before execution events are accepted");
assert_eq!(audit.policy_revision, 7);
assert_eq!(audit.executor_mapping, "test_actor:permission_mode=default");
}
#[tokio::test]
async fn actor_initial_batch_acks_in_order_and_rejects_stale_confirmation() {
let temp = tempfile::tempdir().unwrap();
let store = Arc::new(
bamboo_storage::SessionStoreV2::new(temp.path().to_path_buf())
.await
.unwrap(),
);
let storage: Arc<dyn Storage> = store.clone();
let locked = Arc::new(bamboo_storage::LockedSessionStore::new(storage.clone()));
let inbox: Arc<dyn SessionInboxPort> = Arc::new(bamboo_storage::FileSessionInbox::new(
store.clone(),
bamboo_domain::SessionInboxLimits::default(),
));
let session_id = "actor-confirmation-order";
let run_id = "actor-run-current";
let mut session = Session::new(session_id, "model");
store.save_session(&session).await.unwrap();
for (id, text) in [("actor-first", "first"), ("actor-second", "second")] {
let mut envelope = bamboo_domain::SessionMessageEnvelope::user_input(session_id, text);
envelope.id = bamboo_domain::SessionMessageId::parse(id).unwrap();
inbox.deliver(&envelope).await.unwrap();
}
inbox
.mark_activation_eligible(
session_id,
2,
bamboo_domain::SessionActivationPolicy::InterruptSpecificWait,
)
.await
.unwrap();
let router = SessionActivationRouter::new();
let mut owner_registration = router.register_run(session_id, run_id).await.unwrap();
let binding = SessionInboxRuntimeBinding {
router,
inbox: inbox.clone(),
storage,
persistence: locked,
};
let pairs = claim_canonical_deliveries(&binding, &mut session, run_id, usize::MAX)
.await
.unwrap();
assert_eq!(
pairs
.iter()
.map(|(claim, _)| claim.envelope.id.as_str())
.collect::<Vec<_>>(),
vec!["actor-first", "actor-second"]
);
let seeded = store.load_session(session_id).await.unwrap().unwrap();
assert_eq!(
seeded
.messages
.iter()
.filter(|message| matches!(message.id.as_str(), "actor-first" | "actor-second"))
.map(|message| message.id.as_str())
.collect::<Vec<_>>(),
vec!["actor-first", "actor-second"],
"host context must be durable before actor dispatch"
);
for (claim, _) in &pairs {
assert_eq!(
seeded
.messages
.iter()
.filter(|message| bamboo_domain::is_matching_session_message(
message,
&claim.envelope
))
.count(),
1,
"pre-dispatch host checkpoint must contain exactly one canonical marker for {}",
claim.envelope.id
);
}
assert!(
seeded.session_inbox_admission().is_none_or(|cursor| {
!cursor.contains(&pairs[0].0.envelope.id)
&& !cursor.contains(&pairs[1].0.envelope.id)
}),
"pre-dispatch transcript seeding must not forge worker confirmation"
);
assert_eq!(inbox.inspect(session_id).await.unwrap().claimed, 2);
let claims = pairs
.into_iter()
.map(|(claim, _)| claim)
.collect::<VecDeque<_>>();
let first = claims[0].clone();
let second = claims[1].clone();
let mut stale = admission_confirmation(session_id, &first, "stale-run");
stale.canonical_claim_generation = second.generation;
let mut link = ConfirmationSequenceLink {
frames: VecDeque::from([
ChildFrame::SessionMessageAdmitted {
confirmation: stale,
},
ChildFrame::SessionMessageAdmitted {
confirmation: admission_confirmation(session_id, &first, run_id),
},
ChildFrame::SessionMessageAdmitted {
confirmation: admission_confirmation(session_id, &second, run_id),
},
ChildFrame::Terminal {
status: TerminalStatus::Completed,
result: Some("done".to_string()),
error: None,
transcript: Vec::new(),
},
]),
sent: Vec::new(),
};
let (event_tx, _event_rx) = mpsc::channel(8);
let cancel = CancellationToken::new();
let (_live_tx, mut live_rx) = mpsc::unbounded_channel();
let (_delivery_tx, mut delivery_rx) = mpsc::unbounded_channel();
let result = drive(ActorDriveContext {
client: &mut link,
parent_session_id: "parent",
child_session_id: session_id,
child_attempt: 0,
approval_registry: None,
approval_decider: None,
approval_reviewer: None,
escalation_bridge: None,
event_tx: &event_tx,
cancel_token: &cancel,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: &mut session,
expected_permission_posture: None,
session_inbox_runtime: Some(&binding),
activation_run_id: Some(run_id),
initial_inflight_claims: claims,
first_frame_timeout: Some(Duration::from_secs(1)),
})
.await
.unwrap();
assert_eq!(result.as_deref(), Some("done"));
assert_eq!(
session
.messages
.iter()
.filter(|message| matches!(message.id.as_str(), "actor-first" | "actor-second"))
.map(|message| message.id.as_str())
.collect::<Vec<_>>(),
vec!["actor-first", "actor-second"]
);
let backlog = inbox.inspect(session_id).await.unwrap();
assert_eq!(backlog.pending + backlog.claimed, 0);
let confirmed = store.load_session(session_id).await.unwrap().unwrap();
let cursor = confirmed
.session_inbox_admission()
.expect("exact worker confirmation must checkpoint the admission cursor");
assert!(cursor.contains(&first.envelope.id));
assert!(cursor.contains(&second.envelope.id));
for claim in [&first, &second] {
assert_eq!(
confirmed
.messages
.iter()
.filter(|message| bamboo_domain::is_matching_session_message(
message,
&claim.envelope
))
.count(),
1,
"confirmation must retain one exact canonical marker for {}",
claim.envelope.id
);
}
assert!(inbox
.was_admitted(session_id, &first.envelope.id)
.await
.unwrap());
assert!(inbox
.was_admitted(session_id, &second.envelope.id)
.await
.unwrap());
owner_registration.begin_finalization().await;
owner_registration.finish(2).await.unwrap();
}
#[derive(Default)]
struct RecordingCodexTokenAuthority {
issued_for: std::sync::Mutex<Vec<String>>,
revoked: std::sync::Mutex<Vec<String>>,
}
impl CodexRunTokenAuthority for RecordingCodexTokenAuthority {
fn issue(&self, session_id: &str) -> Result<IssuedCodexRunToken, String> {
self.issued_for
.lock()
.expect("issued fixture lock")
.push(session_id.to_string());
Ok(IssuedCodexRunToken {
token_id: format!("id-{session_id}"),
token: format!("bcx1_secret-{session_id}"),
})
}
fn revoke(&self, token_id: &str) {
self.revoked
.lock()
.expect("revoked fixture lock")
.push(token_id.to_string());
}
}
fn codex_executor(auth_mode: Option<&str>, inherit_user_config: Option<bool>) -> ExecutorSpec {
let bamboo_mode = auth_mode == Some("bamboo")
|| (auth_mode.is_none() && !inherit_user_config.unwrap_or(false));
ExecutorSpec::Codex {
binary: None,
model: None,
mode: None,
sandbox: None,
inherit_user_config,
auth_mode: auth_mode.map(str::to_string),
base_url: bamboo_mode.then(|| "http://127.0.0.1:9562/openai/v1".to_string()),
wire_api: Some("responses".to_string()),
provider_key_ref: None,
forward_env: None,
approval_policy: None,
network_access: None,
allow_danger_bypass: None,
permission_profile: None,
workspace_owned: None,
}
}
fn permission_resolution(
requested: bamboo_domain::SessionPermissionMode,
effective: bamboo_domain::PermissionMode,
) -> bamboo_domain::PermissionModeResolution {
bamboo_domain::PermissionModeResolution {
requested,
effective,
}
}
#[test]
fn permission_posture_mapping_contract_is_exact_for_supported_executors() {
use bamboo_domain::{PermissionMode, SessionPermissionMode};
let default =
permission_resolution(SessionPermissionMode::Default, PermissionMode::Default);
assert_eq!(
expected_permission_executor_mapping(&ExecutorSpec::BambooRuntime, default, false)
.unwrap()
.as_deref(),
Some("bamboo_runtime:default")
);
assert_eq!(
expected_permission_executor_mapping(&ExecutorSpec::Echo, default, false).unwrap(),
None,
"transport-only Echo must not claim the typed permission contract"
);
assert_eq!(
expected_permission_executor_mapping(
&ExecutorSpec::CliAdapter {
command: "must-not-appear-in-contract".to_string(),
args: vec!["credential-like-argument".to_string()],
},
default,
false,
)
.unwrap(),
None,
"unimplemented CliAdapter must not leak command data into a contract"
);
let claude = ExecutorSpec::ClaudeCode {
binary: None,
model: None,
permission_mode: Some("default".to_string()),
inherit_user_config: None,
forward_env: None,
};
for (resolution, mapping) in [
(
permission_resolution(SessionPermissionMode::Default, PermissionMode::Plan),
"claude_code:permission_mode=plan",
),
(
permission_resolution(SessionPermissionMode::Auto, PermissionMode::Auto),
"claude_code:permission_mode=bypassPermissions",
),
(
permission_resolution(SessionPermissionMode::Default, PermissionMode::AcceptEdits),
"claude_code:permission_mode=acceptEdits",
),
(
permission_resolution(SessionPermissionMode::Default, PermissionMode::DontAsk),
"claude_code:permission_mode=dontAsk",
),
(
permission_resolution(
SessionPermissionMode::Bypass,
PermissionMode::BypassPermissions,
),
"claude_code:permission_mode=default",
),
] {
assert_eq!(
expected_permission_executor_mapping(&claude, resolution, false)
.unwrap()
.as_deref(),
Some(mapping)
);
}
assert_eq!(
expected_permission_executor_mapping(&claude, default, true)
.unwrap()
.as_deref(),
Some("claude_code:blocked_explicit_deny")
);
let mut codex_exec = codex_executor(Some("inherit"), Some(true));
if let ExecutorSpec::Codex {
approval_policy, ..
} = &mut codex_exec
{
*approval_policy = Some("on-failure".to_string());
}
assert_eq!(
expected_permission_executor_mapping(&codex_exec, default, false)
.unwrap()
.as_deref(),
Some("codex_exec:approval_policy=on-failure")
);
assert_eq!(
expected_permission_executor_mapping(
&codex_exec,
permission_resolution(SessionPermissionMode::Auto, PermissionMode::Auto),
false,
)
.unwrap()
.as_deref(),
Some("codex_exec:approval_policy=never")
);
assert_eq!(
expected_permission_executor_mapping(&codex_exec, default, true)
.unwrap()
.as_deref(),
Some("codex_exec:blocked_explicit_deny")
);
let mut codex_app_server = codex_executor(Some("inherit"), Some(true));
if let ExecutorSpec::Codex {
mode,
approval_policy,
..
} = &mut codex_app_server
{
*mode = Some("app_server".to_string());
*approval_policy = Some("on-request".to_string());
}
assert_eq!(
expected_permission_executor_mapping(&codex_app_server, default, false)
.unwrap()
.as_deref(),
Some("codex_app_server:approvalPolicy=on-request")
);
assert_eq!(
expected_permission_executor_mapping(
&codex_app_server,
permission_resolution(SessionPermissionMode::Auto, PermissionMode::Auto),
false,
)
.unwrap()
.as_deref(),
Some("codex_app_server:approvalPolicy=never")
);
assert_eq!(
expected_permission_executor_mapping(&codex_app_server, default, true)
.unwrap()
.as_deref(),
Some("codex_app_server:blocked_explicit_deny")
);
}
#[test]
fn only_bamboo_managed_non_git_workspaces_are_marked_owned() {
let project = tempfile::tempdir().unwrap();
let managed = project.path().join(".bamboo/worktree/child-571");
std::fs::create_dir_all(&managed).unwrap();
assert!(!workspace_is_bamboo_owned(managed.to_str().unwrap()));
let marker = project
.path()
.join(".bamboo/worktree/.bamboo-owned/child-571");
std::fs::create_dir_all(marker.parent().unwrap()).unwrap();
std::fs::write(&marker, "bamboo/child-571").unwrap();
assert!(workspace_is_bamboo_owned(managed.to_str().unwrap()));
let nested = managed.join("nested/path");
std::fs::create_dir_all(&nested).unwrap();
assert!(workspace_is_bamboo_owned(nested.to_str().unwrap()));
let arbitrary = tempfile::tempdir().unwrap();
assert!(!workspace_is_bamboo_owned(
arbitrary.path().to_str().unwrap()
));
}
#[test]
fn bamboo_codex_token_is_per_run_redacted_and_revoked_on_guard_drop() {
let authority = Arc::new(RecordingCodexTokenAuthority::default());
let authority_dyn: Arc<dyn CodexRunTokenAuthority> = authority.clone();
let (secrets, guard) = build_codex_run_secrets(
&codex_executor(Some("bamboo"), None),
Some(authority_dyn),
"child-570",
)
.unwrap();
let token = secrets
.codex_provider_token
.as_ref()
.expect("bamboo mode mints a token");
assert_eq!(token.expose(), "bcx1_secret-child-570");
assert!(!format!("{token:?}").contains("secret-child-570"));
assert_eq!(
authority.issued_for.lock().unwrap().as_slice(),
["child-570"]
);
assert!(authority.revoked.lock().unwrap().is_empty());
drop(guard);
assert_eq!(
authority.revoked.lock().unwrap().as_slice(),
["id-child-570"]
);
}
#[test]
fn non_bamboo_codex_never_mints_and_bamboo_fails_closed_without_authority() {
let authority = Arc::new(RecordingCodexTokenAuthority::default());
let authority_dyn: Arc<dyn CodexRunTokenAuthority> = authority.clone();
let (secrets, guard) = build_codex_run_secrets(
&codex_executor(Some("custom"), None),
Some(authority_dyn),
"child-custom",
)
.unwrap();
assert!(secrets.codex_provider_token.is_none());
assert!(guard.is_none());
assert!(authority.issued_for.lock().unwrap().is_empty());
let error = build_codex_run_secrets(
&codex_executor(Some("bamboo"), None),
None,
"child-no-authority",
)
.err()
.expect("bamboo mode without an authority must fail closed");
assert!(error.to_string().contains("per-run token authority"));
}
#[test]
fn codex_provisioning_never_leaks_the_session_provider_credential() {
let credentials = vec![ScopedCredential {
provider: "openai".to_string(),
api_key: "upstream-secret-must-not-cross".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.openai.api_key".to_string()),
}];
for (mode, label) in [
(Some("inherit"), "inherit"),
(Some("api_key"), "api_key"),
(Some("bamboo"), "bamboo"),
(None, "default-bamboo"),
] {
let runner = ActorChildRunner::new(
format!("codex-{label}-test"),
PathBuf::from("/bin/false"),
Vec::new(),
std::env::temp_dir().join(format!("bamboo-codex-{label}-570")),
codex_executor(mode, None),
credentials.clone(),
"openai".to_string(),
1,
);
let mut session = Session::new(format!("child-{label}"), "model");
session.add_message(bamboo_agent_core::Message::user("test"));
let spec = runner.build_spec(
&session,
&crate::runtime::execution::SpawnJob {
parent_session_id: "parent".to_string(),
child_session_id: format!("child-{label}"),
model: "gpt-5.4".to_string(),
disabled_tools: None,
},
);
assert!(
spec.secrets.provider_credentials.is_empty(),
"{label} Codex must not receive the session provider key"
);
}
}
#[test]
fn non_codex_provisioning_still_receives_only_its_selected_provider_credential() {
let credentials = vec![
ScopedCredential {
provider: "openai".to_string(),
api_key: "selected-openai-secret".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.openai.api_key".to_string()),
},
ScopedCredential {
provider: "other".to_string(),
api_key: "unrelated-secret".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.other.api_key".to_string()),
},
];
let runner = ActorChildRunner::new(
"echo-test".to_string(),
PathBuf::from("/bin/false"),
Vec::new(),
std::env::temp_dir().join("bamboo-echo-provider-570"),
ExecutorSpec::Echo,
credentials,
"openai".to_string(),
1,
);
let spec = runner.build_spec(
&Session::new("child-echo", "model"),
&crate::runtime::execution::SpawnJob {
parent_session_id: "parent".to_string(),
child_session_id: "child-echo".to_string(),
model: "gpt-5.4".to_string(),
disabled_tools: None,
},
);
assert_eq!(spec.secrets.provider_credentials.len(), 1);
assert_eq!(
spec.secrets.provider_credentials[0].api_key,
"selected-openai-secret"
);
}
#[test]
fn custom_codex_provisioning_scopes_only_the_referenced_credential() {
let mut executor = codex_executor(Some("custom"), None);
if let ExecutorSpec::Codex {
base_url,
provider_key_ref,
..
} = &mut executor
{
*base_url = Some("https://provider.example/v1".to_string());
*provider_key_ref = Some("provider.custom.api_key".to_string());
}
let credentials = vec![
ScopedCredential {
provider: "openai".to_string(),
api_key: "session-provider-secret".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.openai.api_key".to_string()),
},
ScopedCredential {
provider: "custom".to_string(),
api_key: "selected-secret".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.custom.api_key".to_string()),
},
ScopedCredential {
provider: "other".to_string(),
api_key: "unrelated-secret".to_string(),
base_url: None,
provider_type: Some("openai".to_string()),
credential_ref: Some("provider.other.api_key".to_string()),
},
];
let runner = ActorChildRunner::new(
"codex-test".to_string(),
PathBuf::from("/bin/false"),
Vec::new(),
std::env::temp_dir().join("bamboo-codex-570"),
executor,
credentials,
"openai".to_string(),
1,
);
let mut session = Session::new("child-custom", "model");
session.add_message(bamboo_agent_core::Message::user("test"));
let spec = runner.build_spec(
&session,
&crate::runtime::execution::SpawnJob {
parent_session_id: "parent".to_string(),
child_session_id: "child-custom".to_string(),
model: "gpt-5.4".to_string(),
disabled_tools: None,
},
);
assert_eq!(spec.secrets.provider_credentials.len(), 1);
assert_eq!(
spec.secrets.provider_credentials[0]
.credential_ref
.as_deref(),
Some("provider.custom.api_key")
);
assert_eq!(
spec.secrets.provider_credentials[0].api_key,
"selected-secret"
);
}
fn spec_with(
role: &str,
provider: &str,
model: &str,
workspace: Option<&str>,
disabled: Option<Vec<&str>>,
) -> ProvisionSpec {
let mut spec = ProvisionSpec::new(
ChildIdentity {
child_id: "c".into(),
parent_id: None,
project_key: None,
role: role.into(),
depth: 0,
},
ExecutorSpec::Echo,
"/tmp/fab".into(),
);
spec.workspace = workspace.map(|w| w.to_string());
spec.model = Some(ModelRefSpec {
provider: provider.into(),
model: model.into(),
});
spec.disabled_tools = disabled.map(|d| d.into_iter().map(String::from).collect());
spec
}
#[test]
fn fingerprint_matches_interchangeable_children() {
let a = spec_with(
"explorer",
"p",
"m",
Some("/ws"),
Some(vec!["Bash", "Edit"]),
);
let mut b = spec_with(
"explorer",
"p",
"m",
Some("/ws"),
Some(vec!["Edit", "Bash"]),
);
b.identity.child_id = "other".into();
assert_eq!(
ActorChildRunner::fingerprint(&a),
ActorChildRunner::fingerprint(&b)
);
}
#[test]
fn logical_identity_is_invariant_across_local_remote_scheduled_and_warm_reuse() {
let mut session =
Session::new_child("logical-child-681", "logical-parent-681", "model", "child");
session.root_session_id = "logical-root-681".to_string();
let job = SpawnJob {
parent_session_id: "logical-parent-681".to_string(),
child_session_id: "logical-child-681".to_string(),
model: "model".to_string(),
disabled_tools: None,
};
let expected = LogicalSessionIdentity {
session_id: "logical-child-681".to_string(),
parent_session_id: Some("logical-parent-681".to_string()),
root_session_id: "logical-root-681".to_string(),
};
let placements_and_transport_ids = [
(Placement::Local, "local-mailbox-first"),
(
Placement::Remote {
endpoint: "wss://remote.example/actor".to_string(),
},
"remote-process-44",
),
(
Placement::Schedulable {
pool: "gpu-pool".to_string(),
},
"scheduled-mailbox-9",
),
(Placement::Local, "warm-mailbox-reused-77"),
];
for (placement, transport_id) in placements_and_transport_ids {
let mut provision = spec_with("worker", "provider", "model", None, None);
provision.placement = placement;
provision.identity.child_id = transport_id.to_string();
assert_eq!(logical_identity_for_actor_run(&session, &job), expected);
assert_ne!(
provision.identity.child_id, expected.session_id,
"test fixture must prove transport identity is independent"
);
}
}
#[test]
fn fingerprint_separates_distinct_runtimes() {
let base = spec_with("explorer", "p", "m", Some("/ws"), None);
let base_fp = ActorChildRunner::fingerprint(&base);
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&spec_with("writer", "p", "m", Some("/ws"), None))
);
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&spec_with("explorer", "p2", "m", Some("/ws"), None))
);
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&spec_with("explorer", "p", "m2", Some("/ws"), None))
);
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&spec_with("explorer", "p", "m", Some("/ws2"), None))
);
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&spec_with(
"explorer",
"p",
"m",
Some("/ws"),
Some(vec!["Bash"])
))
);
}
#[test]
fn fingerprint_splits_on_baked_capabilities() {
let base_fp =
ActorChildRunner::fingerprint(&spec_with("explorer", "p", "m", Some("/ws"), None));
let mut depth = spec_with("explorer", "p", "m", Some("/ws"), None);
depth.identity.depth = 2;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&depth),
"depth must split"
);
let mut nested = spec_with("explorer", "p", "m", Some("/ws"), None);
nested.capabilities.nested_spawn = true;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&nested),
"nested_spawn must split"
);
let mut bypass = spec_with("explorer", "p", "m", Some("/ws"), None);
bypass.capabilities.bypass = true;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&bypass),
"bypass must split"
);
let mut auto = spec_with("explorer", "p", "m", Some("/ws"), None);
auto.capabilities.auto_approve_permissions = true;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&auto),
"auto_approve_permissions must split"
);
let mut global_auto = spec_with("explorer", "p", "m", Some("/ws"), None);
global_auto.capabilities.permission_requested_mode = "default".to_string();
global_auto.capabilities.permission_effective_mode = "auto".to_string();
global_auto.capabilities.auto_approve_permissions = true;
let mut explicit_auto = global_auto.clone();
explicit_auto.capabilities.permission_requested_mode = "auto".to_string();
assert_ne!(
ActorChildRunner::fingerprint(&global_auto),
ActorChildRunner::fingerprint(&explicit_auto),
"permission_requested_mode must split global and explicit Auto"
);
let mut plan_overlay = explicit_auto.clone();
plan_overlay.capabilities.permission_effective_mode = "plan".to_string();
assert_ne!(
ActorChildRunner::fingerprint(&explicit_auto),
ActorChildRunner::fingerprint(&plan_overlay),
"permission_effective_mode must split Plan overlay from Auto"
);
let mut enforce = spec_with("explorer", "p", "m", Some("/ws"), None);
enforce.capabilities.enforce_permissions = true;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&enforce),
"enforce_permissions must split"
);
let mut cap = spec_with("explorer", "p", "m", Some("/ws"), None);
cap.capabilities.max_spawn_depth = Some(8);
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&cap),
"max_spawn_depth must split"
);
let mut nha = spec_with("explorer", "p", "m", Some("/ws"), None);
nha.capabilities.no_human_approver = true;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&nha),
"no_human_approver must split"
);
let mut gro = spec_with("explorer", "p", "m", Some("/ws"), None);
gro.capabilities.guardian_read_only = true;
assert_ne!(
base_fp,
ActorChildRunner::fingerprint(&gro),
"guardian_read_only must split"
);
}
#[test]
fn fingerprint_splits_codex_exec_and_app_server_workers() {
let mut exec = spec_with("explorer", "p", "m", Some("/ws"), None);
exec.executor = codex_executor(Some("inherit"), None);
let mut app_server = exec.clone();
if let ExecutorSpec::Codex { mode, .. } = &mut app_server.executor {
*mode = Some("app_server".to_string());
}
assert_ne!(
ActorChildRunner::fingerprint(&exec),
ActorChildRunner::fingerprint(&app_server)
);
}
struct StaticDecider(bool);
#[async_trait]
impl ChildApprovalDecider for StaticDecider {
async fn decide(&self, _child: &str, _req: &serde_json::Value) -> bool {
self.0
}
}
struct RecordingReviewer {
reviewed: mpsc::UnboundedSender<(String, String, serde_json::Value)>,
}
#[async_trait]
impl ChildApprovalReviewer for RecordingReviewer {
async fn review(&self, parent: &str, child: &str, request: &serde_json::Value) -> bool {
let _ = self
.reviewed
.send((parent.to_string(), child.to_string(), request.clone()));
true
}
}
struct SilentLink;
#[async_trait]
impl bamboo_subagent::ChildLink for SilentLink {
async fn send(&mut self, _: ParentFrame) -> bamboo_subagent::TransportResult<()> {
Ok(())
}
async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
std::future::pending().await
}
}
struct InstantTerminalLink {
done: bool,
}
struct ApprovalRoundTripLink {
step: u8,
approval_reply: Option<(String, bool)>,
}
#[async_trait]
impl bamboo_subagent::ChildLink for ApprovalRoundTripLink {
async fn send(&mut self, frame: ParentFrame) -> bamboo_subagent::TransportResult<()> {
if let ParentFrame::ApprovalReply { id, approved } = frame {
self.approval_reply = Some((id, approved));
self.step = 2;
}
Ok(())
}
async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
match self.step {
0 => {
self.step = 1;
Ok(Some(ChildFrame::ApprovalRequest {
id: "approval-1".into(),
body: serde_json::json!({
"tool_name": "Bash",
"permission": "execute",
"resource": "rm -rf target",
"permission_request": {"reason_code": "hard_dangerous"}
}),
}))
}
1 => std::future::pending().await,
2 => {
self.step = 3;
Ok(Some(ChildFrame::Terminal {
status: TerminalStatus::Completed,
result: Some("done".into()),
error: None,
transcript: vec![],
}))
}
_ => std::future::pending().await,
}
}
}
#[tokio::test]
async fn drive_routes_forced_ask_to_parent_reviewer_without_human_event() {
let (event_tx, mut event_rx) = mpsc::channel::<AgentEvent>(8);
let (review_tx, mut review_rx) = mpsc::unbounded_channel();
let reviewer: Arc<dyn ChildApprovalReviewer> = Arc::new(RecordingReviewer {
reviewed: review_tx,
});
let cancel = CancellationToken::new();
let (live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
let (_delivery_tx, mut delivery_rx) = mpsc::unbounded_channel();
let mut logical_session = Session::new("child-reviewer", "model");
let live_guard = crate::external_agents::live::register("child-reviewer", live_tx, 0, None);
let mut link = ApprovalRoundTripLink {
step: 0,
approval_reply: None,
};
let result = tokio::time::timeout(
Duration::from_secs(1),
drive(ActorDriveContext {
client: &mut link,
parent_session_id: "parent-reviewer",
child_session_id: "child-reviewer",
child_attempt: 0,
approval_registry: None,
approval_decider: None,
approval_reviewer: Some(&reviewer),
escalation_bridge: None,
event_tx: &event_tx,
cancel_token: &cancel,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: &mut logical_session,
expected_permission_posture: None,
session_inbox_runtime: None,
activation_run_id: None,
initial_inflight_claims: VecDeque::new(),
first_frame_timeout: None,
}),
)
.await
.expect("worker must receive the reviewer verdict before terminating");
assert_eq!(result.ok().flatten().as_deref(), Some("done"));
assert_eq!(
link.approval_reply,
Some(("approval-1".to_string(), true)),
"reviewer verdict must traverse the live route back to the worker"
);
let (parent, child, body) = tokio::time::timeout(Duration::from_secs(1), review_rx.recv())
.await
.expect("reviewer should be invoked off-loop")
.expect("review channel should remain open");
assert_eq!(parent, "parent-reviewer");
assert_eq!(child, "child-reviewer");
assert_eq!(
body.pointer("/permission_request/reason_code")
.and_then(serde_json::Value::as_str),
Some("hard_dangerous")
);
assert!(
event_rx.try_recv().is_err(),
"must not emit a human-review event"
);
drop(live_guard);
}
#[tokio::test]
async fn drive_denies_forced_ask_without_parent_reviewer_or_manual_event() {
let (event_tx, mut event_rx) = mpsc::channel::<AgentEvent>(8);
let cancel = CancellationToken::new();
let (_live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
let (_delivery_tx, mut delivery_rx) = mpsc::unbounded_channel();
let mut logical_session = Session::new("child-no-reviewer", "model");
let mut link = ApprovalRoundTripLink {
step: 0,
approval_reply: None,
};
let result = tokio::time::timeout(
Duration::from_secs(1),
drive(ActorDriveContext {
client: &mut link,
parent_session_id: "parent-no-reviewer",
child_session_id: "child-no-reviewer",
child_attempt: 0,
approval_registry: None,
approval_decider: None,
approval_reviewer: None,
escalation_bridge: None,
event_tx: &event_tx,
cancel_token: &cancel,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: &mut logical_session,
expected_permission_posture: None,
session_inbox_runtime: None,
activation_run_id: None,
initial_inflight_claims: VecDeque::new(),
first_frame_timeout: None,
}),
)
.await
.expect("fail-closed reply must unblock the child immediately");
assert_eq!(result.ok().flatten().as_deref(), Some("done"));
assert_eq!(link.approval_reply, Some(("approval-1".to_string(), false)));
assert!(
event_rx.try_recv().is_err(),
"missing parent review must not surface a manual approval event"
);
}
#[async_trait]
impl bamboo_subagent::ChildLink for InstantTerminalLink {
async fn send(&mut self, _: ParentFrame) -> bamboo_subagent::TransportResult<()> {
Ok(())
}
async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
if self.done {
std::future::pending().await
} else {
self.done = true;
Ok(Some(ChildFrame::Terminal {
status: TerminalStatus::Completed,
result: Some("done".into()),
error: None,
transcript: vec![],
}))
}
}
}
#[tokio::test]
async fn drive_trips_first_frame_watchdog_on_a_silent_worker() {
let (event_tx, _rx) = mpsc::channel::<AgentEvent>(8);
let cancel = CancellationToken::new();
let (_live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
let (_delivery_tx, mut delivery_rx) = mpsc::unbounded_channel();
let mut logical_session = Session::new("child-x", "model");
let mut link = SilentLink;
let r = drive(ActorDriveContext {
client: &mut link,
parent_session_id: "parent-x",
child_session_id: "child-x",
child_attempt: 0,
approval_registry: None,
approval_decider: None,
approval_reviewer: None,
escalation_bridge: None,
event_tx: &event_tx,
cancel_token: &cancel,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: &mut logical_session,
expected_permission_posture: None,
session_inbox_runtime: None,
activation_run_id: None,
initial_inflight_claims: VecDeque::new(),
first_frame_timeout: Some(Duration::from_millis(100)),
})
.await;
assert!(
matches!(r, Err(AgentError::WorkerUnresponsive(_))),
"a silent worker must trip the first-frame watchdog, got {r:?}"
);
}
#[tokio::test]
async fn drive_does_not_trip_when_a_frame_arrives() {
let (event_tx, _rx) = mpsc::channel::<AgentEvent>(8);
let cancel = CancellationToken::new();
let (_live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
let (_delivery_tx, mut delivery_rx) = mpsc::unbounded_channel();
let mut logical_session = Session::new("child-y", "model");
let mut link = InstantTerminalLink { done: false };
let r = drive(ActorDriveContext {
client: &mut link,
parent_session_id: "parent-y",
child_session_id: "child-y",
child_attempt: 0,
approval_registry: None,
approval_decider: None,
approval_reviewer: None,
escalation_bridge: None,
event_tx: &event_tx,
cancel_token: &cancel,
live_rx: &mut live_rx,
delivery_rx: &mut delivery_rx,
logical_session: &mut logical_session,
expected_permission_posture: None,
session_inbox_runtime: None,
activation_run_id: None,
initial_inflight_claims: VecDeque::new(),
first_frame_timeout: Some(Duration::from_millis(50)),
})
.await;
assert_eq!(r.ok().flatten().as_deref(), Some("done"));
}
#[tokio::test]
async fn child_approval_fails_closed_without_decider() {
let body = serde_json::json!({"tool_name":"Bash","permission":"run","resource":"rm -rf /"});
assert!(!decide_child_approval(None, "child-1", &body).await);
}
#[tokio::test]
async fn child_approval_honors_wired_decider() {
let body =
serde_json::json!({"tool_name":"Write","permission":"write","resource":"/tmp/x"});
let approve: Arc<dyn ChildApprovalDecider> = Arc::new(StaticDecider(true));
let deny: Arc<dyn ChildApprovalDecider> = Arc::new(StaticDecider(false));
assert!(decide_child_approval(Some(&approve), "child-1", &body).await);
assert!(!decide_child_approval(Some(&deny), "child-1", &body).await);
}
use crate::runtime::execution::SpawnJob;
use bamboo_agent_core::Session;
fn bogus_runner(placements: HashMap<String, ResolvedRemotePlacement>) -> ActorChildRunner {
ActorChildRunner::new(
"test-actor".into(),
PathBuf::from("/bin/false"),
vec![],
std::env::temp_dir().join("bamboo-test-fab-193"),
ExecutorSpec::Echo,
vec![],
"anthropic".into(),
4,
)
.with_remote_placements(placements)
}
fn session_of_role(role: &str, assignment: &str) -> Session {
let mut s = Session::new("child-1", "test-model");
s.metadata
.insert("subagent_type".to_string(), role.to_string());
s.add_message(bamboo_agent_core::Message::user(assignment));
s
}
fn job_for(child: &str) -> SpawnJob {
SpawnJob {
parent_session_id: "parent-1".into(),
child_session_id: child.into(),
model: String::new(),
disabled_tools: None,
}
}
#[derive(Default)]
struct RecordingChildSessionPort {
saved: std::sync::Mutex<Option<Session>>,
}
impl RecordingChildSessionPort {
fn saved_child(&self) -> Session {
self.saved
.lock()
.expect("saved-child fixture lock")
.clone()
.expect("create_child_action must save the child")
}
}
#[async_trait]
impl crate::session_app::child_session::ChildSessionPort for RecordingChildSessionPort {
async fn load_root_session(
&self,
_root_id: &str,
) -> Result<Session, crate::session_app::child_session::ChildSessionError> {
unreachable!("create_child_action does not load the root")
}
async fn load_child_for_parent(
&self,
_parent_id: &str,
_child_id: &str,
) -> Result<Session, crate::session_app::child_session::ChildSessionError> {
unreachable!("create_child_action does not reload the child")
}
async fn save_child_session(
&self,
child: &mut Session,
) -> Result<(), crate::session_app::child_session::ChildSessionError> {
*self.saved.lock().expect("saved-child fixture lock") = Some(child.clone());
Ok(())
}
async fn save_child_session_authoritative_flags(
&self,
_child: &mut Session,
) -> Result<(), crate::session_app::child_session::ChildSessionError> {
unreachable!("new-child creation uses the ordinary save")
}
async fn is_child_running(&self, _child_id: &str) -> bool {
false
}
async fn list_children(
&self,
_parent_id: &str,
) -> Vec<crate::session_app::child_session::ChildSessionEntry> {
Vec::new()
}
async fn enqueue_child_run(
&self,
_parent: &Session,
_child: &Session,
) -> Result<(), crate::session_app::child_session::ChildSessionError> {
unreachable!("fixture creates the child with auto_run=false")
}
async fn cancel_child_run_and_wait(
&self,
_child_id: &str,
) -> Result<(), crate::session_app::child_session::ChildSessionError> {
unreachable!("create_child_action does not cancel")
}
async fn delete_child_session(
&self,
_parent_id: &str,
_child_id: &str,
) -> Result<
crate::session_app::child_session::DeleteChildResult,
crate::session_app::child_session::ChildSessionError,
> {
unreachable!("create_child_action does not delete")
}
async fn get_child_runner_info(
&self,
_child_id: &str,
) -> Option<crate::session_app::child_session::ChildRunnerInfo> {
None
}
async fn register_parent_wait_for_child(
&self,
_parent_session_id: &str,
_child_session_id: &str,
_tool_call_id: Option<&str>,
) -> Result<(), crate::session_app::child_session::ChildSessionError> {
unreachable!("create_child_action does not register a wait")
}
async fn register_parent_wait_for_children(
&self,
_parent_session_id: &str,
_child_session_ids: &[String],
_policy: bamboo_domain::session::runtime_state::ChildWaitPolicy,
) -> Result<usize, crate::session_app::child_session::ChildSessionError> {
unreachable!("create_child_action does not register a wait")
}
async fn active_child_ids(&self, _parent_session_id: &str) -> Vec<String> {
Vec::new()
}
async fn find_resident_child(
&self,
_root_session_id: &str,
_resident_name: &str,
) -> Option<String> {
None
}
async fn ensure_child_indexed(&self, _child_session_id: &str) {}
}
#[test]
fn build_spec_sets_remote_placement_for_matching_role() {
let mut placements = HashMap::new();
placements.insert(
"explorer".to_string(),
ResolvedRemotePlacement {
endpoint: "wss://gpu-host:8443".into(),
token: Some("T-secret".into()),
ca_cert_file: None,
host_label: None,
},
);
let runner = bogus_runner(placements);
let s = session_of_role("explorer", "do the thing");
let spec = runner.build_spec(&s, &job_for("child-1"));
match &spec.placement {
Placement::Remote { endpoint } => assert_eq!(endpoint, "wss://gpu-host:8443"),
other => panic!("expected Remote, got {other:?}"),
}
assert_eq!(spec.secrets.worker_auth_token.as_deref(), Some("T-secret"));
}
#[test]
fn build_spec_leaves_local_for_unmatched_role() {
let mut placements = HashMap::new();
placements.insert(
"explorer".to_string(),
ResolvedRemotePlacement {
endpoint: "wss://gpu-host:8443".into(),
token: Some("T".into()),
ca_cert_file: None,
host_label: None,
},
);
let runner = bogus_runner(placements);
let s = session_of_role("writer", "do the thing");
let spec = runner.build_spec(&s, &job_for("child-1"));
assert_eq!(spec.placement, Placement::Local);
assert!(spec.secrets.worker_auth_token.is_none());
}
#[test]
fn build_spec_local_when_no_placements() {
let runner = bogus_runner(HashMap::new());
let s = session_of_role("explorer", "do the thing");
let spec = runner.build_spec(&s, &job_for("child-1"));
assert_eq!(spec.placement, Placement::Local);
assert!(spec.secrets.worker_auth_token.is_none());
}
#[tokio::test]
async fn build_spec_preserves_exact_inherited_permission_mode_for_child_worker() {
for (label, mode) in [
("bypass", bamboo_domain::SessionPermissionMode::Bypass),
("auto", bamboo_domain::SessionPermissionMode::Auto),
] {
let runner = bogus_runner(HashMap::new());
let mut parent = Session::new(format!("parent-{label}"), "test-model");
parent
.agent_runtime_state
.get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
.set_permission_mode(mode);
let workspace = tempfile::tempdir().expect("workspace fixture");
let port = RecordingChildSessionPort::default();
let child_id = format!("child-{label}-{}", uuid::Uuid::new_v4());
crate::session_app::child_session::create_child_action(
&port,
crate::session_app::child_session::CreateChildInput {
parent_session: parent,
child_id: child_id.clone(),
title: format!("{label} child"),
responsibility: "Run ordinary commands".to_string(),
assignment_prompt: "run an ordinary command".to_string(),
subagent_type: "explorer".to_string(),
workspace: workspace.path().to_string_lossy().into_owned(),
workspace_source: crate::project_context::WorkspaceSource::Explicit,
model_override: None,
model_ref_override: None,
runtime_metadata: HashMap::new(),
auto_run: false,
reasoning_effort: None,
lifecycle: None,
resident_name: None,
resident_context: None,
disabled_tools: None,
context_fork: None,
},
)
.await
.expect("create inherited-permission child");
let child = port.saved_child();
assert_eq!(
child
.agent_runtime_state
.as_ref()
.map(bamboo_domain::AgentRuntimeState::effective_permission_mode),
Some(mode),
"create_child_action must inherit {label} from the parent"
);
let spec = runner.build_spec(&child, &job_for(&child_id));
assert_eq!(
spec.capabilities.bypass,
mode == bamboo_domain::SessionPermissionMode::Bypass
);
assert_eq!(
spec.capabilities.auto_approve_permissions,
mode == bamboo_domain::SessionPermissionMode::Auto
);
assert!(
spec.capabilities.enforce_permissions,
"policy evaluation must remain active under {label}"
);
}
}
#[tokio::test]
async fn child_resident_and_guardian_inherit_project_through_actor_run_spec() {
let project_id = bamboo_domain::ProjectId::parse("project-inherited").expect("Project id");
let workspace = tempfile::tempdir().expect("workspace fixture");
for (role, lifecycle, resident_name) in [
("explorer", None, None),
("resident", Some("resident"), Some("stable-reviewer")),
("guardian", None, None),
] {
let mut parent = Session::new(format!("parent-{role}"), "test-model");
parent.set_project_id_meta(project_id.to_string());
let port = RecordingChildSessionPort::default();
let child_id = format!("child-{role}-{}", uuid::Uuid::new_v4());
crate::session_app::child_session::create_child_action(
&port,
crate::session_app::child_session::CreateChildInput {
parent_session: parent,
child_id: child_id.clone(),
title: format!("{role} child"),
responsibility: "Review the assigned work".to_string(),
assignment_prompt: "inspect the change".to_string(),
subagent_type: role.to_string(),
workspace: workspace.path().to_string_lossy().into_owned(),
workspace_source: crate::project_context::WorkspaceSource::Explicit,
model_override: None,
model_ref_override: None,
runtime_metadata: HashMap::new(),
auto_run: false,
reasoning_effort: None,
lifecycle: lifecycle.map(str::to_string),
resident_name: resident_name.map(str::to_string),
resident_context: None,
disabled_tools: None,
context_fork: None,
},
)
.await
.expect("create Project-inheriting child");
let child = port.saved_child();
assert_eq!(
crate::project_context::ProjectContextResolver::project_id_from_session(&child),
Some(project_id.clone()),
"{role} child must inherit its parent's Project"
);
assert_eq!(
project_id_for_actor_run(&child).expect("valid actor Project identity"),
Some(project_id.clone()),
"{role} actor RunSpec must preserve inherited Project identity"
);
}
}
#[test]
fn placement_metadata_stamps_remote_and_schedulable_not_local() {
assert_eq!(placement_metadata(&Placement::Local, None), None);
let r = placement_metadata(
&Placement::Remote {
endpoint: "wss://10.0.0.5:8443/stream".into(),
},
None,
)
.unwrap();
assert!(r.contains(r#""kind":"remote""#), "{r}");
assert!(r.contains(r#""host":"10.0.0.5""#), "{r}");
let labeled = placement_metadata(
&Placement::Remote {
endpoint: "ws://169.254.230.101:8899".into(),
},
Some("mini"),
)
.unwrap();
assert!(labeled.contains(r#""host":"mini""#), "{labeled}");
let s = placement_metadata(
&Placement::Schedulable {
pool: "explorers".into(),
},
Some("mini"),
)
.unwrap();
assert!(s.contains(r#""kind":"remote""#), "{s}");
assert!(s.contains(r#""host":"mini""#), "{s}");
let p: bamboo_storage::SessionPlacement = serde_json::from_str(&labeled).unwrap();
assert_eq!(p.kind, "remote");
assert_eq!(p.host, "mini");
}
#[tokio::test]
async fn execute_external_child_routes_role_to_remote_worker_without_spawning() {
let token = "remote-test-token";
let server = bamboo_subagent::transport::WsServer::bind_with_token(
(std::net::Ipv4Addr::LOCALHOST, 0).into(),
Some(token.to_string()),
)
.await
.expect("bind resident worker");
let endpoint = server.ws_endpoint(); let srv = tokio::spawn(async move {
let _ = server
.serve(Arc::new(bamboo_subagent::executor::EchoExecutor))
.await;
});
let mut placements = HashMap::new();
placements.insert(
"explorer".to_string(),
ResolvedRemotePlacement {
endpoint: endpoint.clone(),
token: Some(token.to_string()),
ca_cert_file: None,
host_label: Some("mini-e2e".into()), },
);
let runner = bogus_runner(placements);
let mut session = session_of_role("explorer", "hello remote");
let job = job_for("child-1");
let (event_tx, mut event_rx) = mpsc::channel::<AgentEvent>(64);
let cancel = CancellationToken::new();
let result = tokio::time::timeout(
Duration::from_secs(10),
runner.execute_external_child(&mut session, &job, event_tx, cancel),
)
.await
.expect("run did not hang")
.expect("remote run succeeded (connected to resident worker, did not spawn)");
let _ = result;
let last = session
.messages
.iter()
.rev()
.find(|m| matches!(m.role, Role::Assistant))
.expect("an assistant reply was written back");
assert!(
last.content.contains("echo:"),
"expected echo reply, got {:?}",
last.content
);
let placement = session
.metadata
.get("placement")
.expect("remote child session stamped with a placement");
assert!(placement.contains(r#""kind":"remote""#), "{placement}");
assert!(placement.contains(r#""host":"mini-e2e""#), "{placement}");
let mut saw_event = false;
while let Ok(Some(_ev)) =
tokio::time::timeout(Duration::from_millis(50), event_rx.recv()).await
{
saw_event = true;
}
let _ = saw_event;
srv.abort();
}
fn bogus_sched_runner(
remote: HashMap<String, ResolvedRemotePlacement>,
sched: HashMap<String, ResolvedSchedulablePlacement>,
) -> ActorChildRunner {
ActorChildRunner::new(
"test-actor".into(),
PathBuf::from("/bin/false"),
vec![],
std::env::temp_dir().join("bamboo-test-fab-181"),
ExecutorSpec::Echo,
vec![],
"anthropic".into(),
4,
)
.with_remote_placements(remote)
.with_schedulable_placements(sched)
}
fn sched_placement(
pool: &str,
_registry_url: impl Into<String>,
) -> ResolvedSchedulablePlacement {
ResolvedSchedulablePlacement {
pool: pool.into(),
host_label: None,
}
}
#[test]
fn build_spec_sets_schedulable_placement_for_matching_role() {
let mut sched = HashMap::new();
sched.insert(
"explorer".to_string(),
sched_placement("gpu-pool", "unused"),
);
let runner = bogus_sched_runner(HashMap::new(), sched);
let s = session_of_role("explorer", "do the thing");
let spec = runner.build_spec(&s, &job_for("child-1"));
match &spec.placement {
Placement::Schedulable { pool } => assert_eq!(pool, "gpu-pool"),
other => panic!("expected Schedulable, got {other:?}"),
}
assert!(spec.secrets.worker_auth_token.is_none());
}
#[test]
fn build_spec_remote_wins_when_role_in_both_maps() {
let mut remote = HashMap::new();
remote.insert(
"explorer".to_string(),
ResolvedRemotePlacement {
endpoint: "wss://fixed-host:8443".into(),
token: Some("T-remote".into()),
ca_cert_file: None,
host_label: None,
},
);
let mut sched = HashMap::new();
sched.insert(
"explorer".to_string(),
sched_placement("gpu-pool", "https://control-plane:9562"),
);
let runner = bogus_sched_runner(remote, sched);
let s = session_of_role("explorer", "do the thing");
let spec = runner.build_spec(&s, &job_for("child-1"));
match &spec.placement {
Placement::Remote { endpoint } => assert_eq!(endpoint, "wss://fixed-host:8443"),
other => panic!("expected Remote (precedence), got {other:?}"),
}
assert_eq!(spec.secrets.worker_auth_token.as_deref(), Some("T-remote"));
}
#[test]
fn build_spec_local_for_unmatched_schedulable_role() {
let mut sched = HashMap::new();
sched.insert(
"explorer".to_string(),
sched_placement("gpu-pool", "https://control-plane:9562"),
);
let runner = bogus_sched_runner(HashMap::new(), sched);
let s = session_of_role("writer", "do the thing");
let spec = runner.build_spec(&s, &job_for("child-1"));
assert_eq!(spec.placement, Placement::Local);
assert!(spec.secrets.worker_auth_token.is_none());
}
#[test]
fn placement_stamp_uses_node_label_for_remote_and_schedulable() {
let mut remote = HashMap::new();
remote.insert(
"explorer".to_string(),
ResolvedRemotePlacement {
endpoint: "ws://169.254.230.101:8899".into(),
token: None,
ca_cert_file: None,
host_label: Some("mini".into()),
},
);
let runner = bogus_runner(remote);
let spec = runner.build_spec(&session_of_role("explorer", "go"), &job_for("c1"));
let stamp = runner
.placement_stamp_for(&spec)
.expect("remote child is stamped");
assert!(stamp.contains(r#""kind":"remote""#), "{stamp}");
assert!(stamp.contains(r#""host":"mini""#), "{stamp}");
let mut remote_nolabel = HashMap::new();
remote_nolabel.insert(
"explorer".to_string(),
ResolvedRemotePlacement {
endpoint: "ws://169.254.230.101:8899".into(),
token: None,
ca_cert_file: None,
host_label: None,
},
);
let r2 = bogus_runner(remote_nolabel);
let spec2 = r2.build_spec(&session_of_role("explorer", "go"), &job_for("c1"));
assert!(r2
.placement_stamp_for(&spec2)
.unwrap()
.contains(r#""host":"169.254.230.101""#));
let mut sched = HashMap::new();
sched.insert(
"mac-mini-monitor".to_string(),
ResolvedSchedulablePlacement {
pool: "mac-mini-monitor".into(),
host_label: Some("mini".into()),
},
);
let sr = bogus_sched_runner(HashMap::new(), sched);
let spec3 = sr.build_spec(&session_of_role("mac-mini-monitor", "go"), &job_for("c1"));
let stamp3 = sr
.placement_stamp_for(&spec3)
.expect("scheduled child is stamped");
assert!(stamp3.contains(r#""kind":"remote""#), "{stamp3}");
assert!(stamp3.contains(r#""host":"mini""#), "{stamp3}");
let local = bogus_runner(HashMap::new());
let spec4 = local.build_spec(&session_of_role("writer", "go"), &job_for("c1"));
assert_eq!(local.placement_stamp_for(&spec4), None);
}
async fn start_bus() -> (String, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let core = std::sync::Arc::new(bamboo_broker::BrokerCore::new(dir.path()));
let server = std::sync::Arc::new(bamboo_broker::BrokerServer::new(core, "t"));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
let _ = server.serve(listener).await;
});
(format!("ws://{addr}"), dir)
}
async fn join_pool(endpoint: &str, id: &str, pool: &str) -> bamboo_broker::BrokerClient {
let mut c = bamboo_broker::BrokerClient::connect(
endpoint,
bamboo_subagent::AgentRef {
session_id: id.into(),
role: Some(pool.into()),
},
"t",
)
.await
.unwrap();
c.subscribe().await.unwrap();
c
}
fn sched_runner_on_bus(endpoint: &str, child_role: &str, pool: &str) -> ActorChildRunner {
let mut sched = HashMap::new();
sched.insert(child_role.to_string(), sched_placement(pool, "unused"));
bogus_sched_runner(HashMap::new(), sched).with_bus(Some(bamboo_subagent::BusEndpoint {
endpoint: endpoint.into(),
token: "t".into(),
}))
}
#[tokio::test]
async fn resolve_schedulable_picks_a_live_bus_worker() {
let (endpoint, _dir) = start_bus().await;
let _w = join_pool(&endpoint, "w-gpu", "gpu-pool").await;
let runner = sched_runner_on_bus(&endpoint, "explorer", "gpu-pool");
let mailbox = runner
.resolve_schedulable_worker("explorer")
.await
.expect("a live pool worker is found on the bus");
assert_eq!(mailbox, "w-gpu");
}
#[tokio::test]
async fn resolve_schedulable_round_robins_over_pool_workers() {
let (endpoint, _dir) = start_bus().await;
let _a = join_pool(&endpoint, "w-a", "gpu-pool").await;
let _b = join_pool(&endpoint, "w-b", "gpu-pool").await;
let runner = sched_runner_on_bus(&endpoint, "explorer", "gpu-pool");
let mut picked = std::collections::HashSet::new();
for _ in 0..6 {
picked.insert(runner.resolve_schedulable_worker("explorer").await.unwrap());
}
assert_eq!(
picked,
["w-a".to_string(), "w-b".to_string()].into_iter().collect(),
"round-robin must cover every connected pool worker"
);
}
#[tokio::test]
async fn resolve_schedulable_errors_on_empty_pool() {
let (endpoint, _dir) = start_bus().await;
let runner = sched_runner_on_bus(&endpoint, "explorer", "gpu-pool");
let err = runner
.resolve_schedulable_worker("explorer")
.await
.expect_err("an empty pool is terminal — no local fallback")
.to_string();
assert!(err.contains("no live worker in pool"), "got: {err}");
assert!(err.contains("NOT spawning"), "got: {err}");
}
#[tokio::test]
async fn execute_external_child_runs_schedulable_over_bus_and_stamps_node_label() {
let (endpoint, _dir) = start_bus().await;
let ep = endpoint.clone();
let worker = tokio::spawn(async move {
let _ = bamboo_broker::serve_executor(
&ep,
bamboo_subagent::AgentRef {
session_id: "mmm-worker".into(),
role: Some("mac-mini-monitor".into()),
},
"t",
std::sync::Arc::new(bamboo_subagent::executor::EchoExecutor),
)
.await;
});
let mut probe = bamboo_broker::BrokerClient::connect(
&endpoint,
bamboo_subagent::AgentRef {
session_id: "probe".into(),
role: None,
},
"t",
)
.await
.unwrap();
let mut ready = false;
for _ in 0..100 {
if probe
.list_connected("mac-mini-monitor")
.await
.unwrap()
.iter()
.any(|id| id == "mmm-worker")
{
ready = true;
break;
}
tokio::time::sleep(Duration::from_millis(30)).await;
}
assert!(ready, "worker never joined the pool");
let mut sched = HashMap::new();
sched.insert(
"mac-mini-monitor".to_string(),
ResolvedSchedulablePlacement {
pool: "mac-mini-monitor".into(),
host_label: Some("mini".into()),
},
);
let runner = bogus_sched_runner(HashMap::new(), sched).with_bus(Some(
bamboo_subagent::BusEndpoint {
endpoint: endpoint.clone(),
token: "t".into(),
},
));
let mut session = session_of_role("mac-mini-monitor", "hello scheduled");
let job = job_for("child-1");
let (event_tx, _rx) = mpsc::channel::<AgentEvent>(64);
let cancel = CancellationToken::new();
tokio::time::timeout(
Duration::from_secs(10),
runner.execute_external_child(&mut session, &job, event_tx, cancel),
)
.await
.expect("run did not hang")
.expect("schedulable run succeeded over the bus (no local spawn)");
let last = session
.messages
.iter()
.rev()
.find(|m| matches!(m.role, Role::Assistant))
.expect("an assistant reply was written back");
assert!(last.content.contains("echo:"), "got {:?}", last.content);
let placement = session
.metadata
.get("placement")
.expect("scheduled child session stamped with a placement");
assert!(placement.contains(r#""kind":"remote""#), "{placement}");
assert!(placement.contains(r#""host":"mini""#), "{placement}");
worker.abort();
}
}