1use std::collections::HashMap;
14use std::path::PathBuf;
15use std::sync::Arc;
16use std::time::Duration;
17
18use async_trait::async_trait;
19use bamboo_agent_core::{AgentError, AgentEvent, Role, Session};
20use bamboo_domain::poison::PoisonRecover;
21use tokio::sync::mpsc;
22use tokio_util::sync::CancellationToken;
23
24use bamboo_subagent::fleet::{spawn_worker_on_bus, SpawnedChild};
25use bamboo_subagent::proto::{
26 AgentRecord, ChildFrame, ParentFrame, PermissionPolicyContext, RunSpec, TerminalStatus,
27};
28use bamboo_subagent::provision::{
29 ChildIdentity, ExecutorSpec, ModelRefSpec, Placement, ProvisionSpec, ScopedCredential,
30};
31use bamboo_subagent::transport::{client_config_trusting_cert, ChildClient};
32
33use crate::runtime::execution::{ExternalChildRunner, SpawnJob};
34
35pub const DEFAULT_MAX_CONCURRENT_ACTORS: usize = 8;
37
38pub const MAX_SPAWN_DEPTH: u32 = 4;
44
45const DEFAULT_MAX_IDLE_PER_KEY: usize = 4;
47
48const POOLED_IDLE_TIMEOUT_SECS: u64 = 300;
51
52const WORKER_FIRST_FRAME_TIMEOUT: Duration = Duration::from_secs(60);
59
60struct PooledWorker {
66 worker: SpawnedChild,
67 mailbox_id: String,
69}
70
71#[derive(Debug, Clone)]
77pub struct ResolvedRemotePlacement {
78 pub endpoint: String,
79 pub token: Option<String>,
80 pub ca_cert_file: Option<PathBuf>,
81 pub host_label: Option<String>,
85}
86
87#[derive(Debug, Clone)]
94pub struct ResolvedSchedulablePlacement {
95 pub pool: String,
96 pub host_label: Option<String>,
100}
101
102enum PlacementKind {
108 Local,
109 Remote,
110 Schedulable,
111}
112
113pub struct ActorChildRunner {
115 approval_registry: Option<super::approval_registry::SharedApprovalRegistry>,
116 permission_config: Option<Arc<bamboo_tools::permission::PermissionConfig>>,
117 agent_id: String,
118 worker_bin: PathBuf,
119 worker_args: Vec<String>,
120 fabric_dir: PathBuf,
121 executor: ExecutorSpec,
122 credentials: Vec<ScopedCredential>,
125 default_provider: String,
127 bus: Option<bamboo_subagent::BusEndpoint>,
130 concurrency: std::sync::Arc<tokio::sync::Semaphore>,
134 pool: Arc<tokio::sync::Mutex<HashMap<String, Vec<PooledWorker>>>>,
140 max_idle_per_key: usize,
141 approval_decider: Option<Arc<dyn ChildApprovalDecider>>,
145 approval_reviewer: Option<Arc<dyn ChildApprovalReviewer>>,
149 escalation_bridge: Arc<std::sync::Mutex<Option<bamboo_subagent::executor::HostBridge>>>,
160 remote_placements: HashMap<String, ResolvedRemotePlacement>,
166 schedulable_placements: HashMap<String, ResolvedSchedulablePlacement>,
174 schedule_cursor: Arc<std::sync::Mutex<HashMap<String, usize>>>,
179}
180
181#[async_trait]
191pub trait ChildApprovalDecider: Send + Sync {
192 async fn decide(&self, child_session_id: &str, request: &serde_json::Value) -> bool;
195}
196
197async fn decide_child_approval(
200 decider: Option<&Arc<dyn ChildApprovalDecider>>,
201 child_session_id: &str,
202 request: &serde_json::Value,
203) -> bool {
204 match decider {
205 Some(decider) => decider.decide(child_session_id, request).await,
206 None => false,
207 }
208}
209
210const CHILD_APPROVAL_TIMEOUT: Duration = Duration::from_secs(300);
214
215#[async_trait]
224pub trait ChildApprovalReviewer: Send + Sync {
225 async fn review(
228 &self,
229 parent_session_id: &str,
230 child_session_id: &str,
231 request: &serde_json::Value,
232 ) -> bool;
233}
234
235fn child_approval_reviewer_slot() -> &'static std::sync::OnceLock<Arc<dyn ChildApprovalReviewer>> {
236 static SLOT: std::sync::OnceLock<Arc<dyn ChildApprovalReviewer>> = std::sync::OnceLock::new();
237 &SLOT
238}
239
240pub fn set_child_approval_reviewer(reviewer: Arc<dyn ChildApprovalReviewer>) {
242 let _ = child_approval_reviewer_slot().set(reviewer);
243}
244
245pub fn child_approval_reviewer() -> Option<Arc<dyn ChildApprovalReviewer>> {
247 child_approval_reviewer_slot().get().cloned()
248}
249
250impl ActorChildRunner {
251 #[allow(clippy::too_many_arguments)]
252 pub fn new(
253 agent_id: String,
254 worker_bin: PathBuf,
255 worker_args: Vec<String>,
256 fabric_dir: PathBuf,
257 executor: ExecutorSpec,
258 credentials: Vec<ScopedCredential>,
259 default_provider: String,
260 max_concurrent: usize,
261 ) -> Self {
262 Self {
263 approval_registry: None,
264 permission_config: None,
265 agent_id,
266 worker_bin,
267 worker_args,
268 fabric_dir,
269 executor,
270 credentials,
271 default_provider,
272 bus: None,
273 concurrency: std::sync::Arc::new(tokio::sync::Semaphore::new(max_concurrent.max(1))),
274 pool: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
275 max_idle_per_key: DEFAULT_MAX_IDLE_PER_KEY,
276 approval_decider: None,
277 approval_reviewer: None,
278 escalation_bridge: Arc::new(std::sync::Mutex::new(None)),
279 remote_placements: HashMap::new(),
280 schedulable_placements: HashMap::new(),
281 schedule_cursor: Arc::new(std::sync::Mutex::new(HashMap::new())),
282 }
283 }
284
285 pub fn with_approval_registry(
286 mut self,
287 registry: super::approval_registry::SharedApprovalRegistry,
288 ) -> Self {
289 self.approval_registry = Some(registry);
290 self
291 }
292
293 pub fn with_permission_config(
294 mut self,
295 config: Arc<bamboo_tools::permission::PermissionConfig>,
296 ) -> Self {
297 self.permission_config = Some(config);
298 self
299 }
300
301 pub fn with_bus(mut self, bus: Option<bamboo_subagent::BusEndpoint>) -> Self {
306 self.bus = bus.filter(|b| !b.endpoint.trim().is_empty());
307 self
308 }
309
310 pub fn with_approval_decider(mut self, decider: Arc<dyn ChildApprovalDecider>) -> Self {
313 self.approval_decider = Some(decider);
314 self
315 }
316
317 pub fn with_approval_reviewer(mut self, reviewer: Arc<dyn ChildApprovalReviewer>) -> Self {
318 self.approval_reviewer = Some(reviewer);
319 self
320 }
321
322 pub fn with_remote_placements(
328 mut self,
329 placements: HashMap<String, ResolvedRemotePlacement>,
330 ) -> Self {
331 self.remote_placements = placements;
332 self
333 }
334
335 pub fn with_schedulable_placements(
341 mut self,
342 placements: HashMap<String, ResolvedSchedulablePlacement>,
343 ) -> Self {
344 self.schedulable_placements = placements;
345 self
346 }
347
348 fn fingerprint(spec: &ProvisionSpec) -> String {
365 let role = spec.identity.role.as_str();
366 let (provider, model) = spec
367 .model
368 .as_ref()
369 .map(|m| (m.provider.as_str(), m.model.as_str()))
370 .unwrap_or(("", ""));
371 let workspace = spec.workspace.as_deref().unwrap_or("");
372 let mut tools = spec.disabled_tools.clone().unwrap_or_default();
373 tools.sort();
374 let caps = &spec.capabilities;
375 format!(
376 "{role}\u{1}{provider}\u{1}{model}\u{1}{workspace}\u{1}{}\u{1}d={}\u{1}ns={}\u{1}by={}\u{1}ep={}\u{1}md={}\u{1}nha={}\u{1}gro={}",
377 tools.join(","),
378 spec.identity.depth,
379 caps.nested_spawn,
380 caps.bypass,
381 caps.enforce_permissions,
382 caps.max_spawn_depth.unwrap_or(0),
383 caps.no_human_approver,
390 caps.guardian_read_only,
394 )
395 }
396
397 async fn acquire_bus_worker(
403 &self,
404 key: &str,
405 spec: &ProvisionSpec,
406 ) -> crate::runtime::runner::Result<PooledWorker> {
407 loop {
410 let candidate = {
411 let mut pool = self.pool.lock().await;
412 pool.get_mut(key).and_then(|bucket| bucket.pop())
413 };
414 let Some(mut candidate) = candidate else {
415 break;
416 };
417 if candidate.worker.is_alive() {
418 return Ok(candidate);
419 }
420 candidate.worker.kill().await;
421 }
422
423 let spawned = spawn_worker_on_bus(&self.worker_bin, &self.worker_args, spec)
424 .await
425 .map_err(|e| AgentError::LLM(format!("actor spawn (bus) failed: {e}")))?;
426 let mailbox_id = spawned.record.agent_id.clone();
427 Ok(PooledWorker {
428 worker: spawned,
429 mailbox_id,
430 })
431 }
432
433 async fn release_bus_worker(&self, key: &str, mut worker: PooledWorker) {
437 if !worker.worker.is_alive() {
438 worker.worker.kill().await;
439 return;
440 }
441 let mut pool = self.pool.lock().await;
442 let bucket = pool.entry(key.to_string()).or_default();
443 if bucket.len() >= self.max_idle_per_key {
444 drop(pool);
445 worker.worker.kill().await;
446 return;
447 }
448 bucket.push(worker);
449 }
450
451 fn build_spec(&self, session: &Session, job: &SpawnJob) -> ProvisionSpec {
453 let mut spec = ProvisionSpec::new(
454 ChildIdentity {
455 child_id: job.child_session_id.clone(),
456 parent_id: Some(job.parent_session_id.clone()),
457 project_key: None,
458 role: session
459 .metadata
460 .get("subagent_type")
461 .cloned()
462 .unwrap_or_else(|| "worker".to_string()),
463 depth: session.spawn_depth,
468 },
469 self.executor.clone(),
470 self.fabric_dir.to_string_lossy().into_owned(),
471 );
472 spec.workspace = session.workspace.clone();
473 spec.bus = self.bus.clone();
476 spec.model = session
479 .model_ref
480 .as_ref()
481 .map(|r| ModelRefSpec {
482 provider: r.provider.clone(),
483 model: r.model.clone(),
484 })
485 .or_else(|| {
486 let m = job.model.trim();
487 (!m.is_empty()).then(|| ModelRefSpec {
488 provider: self.default_provider.clone(),
489 model: m.to_string(),
490 })
491 });
492 spec.disabled_tools = job.disabled_tools.clone();
493 let provider = spec
495 .model
496 .as_ref()
497 .map(|m| m.provider.as_str())
498 .filter(|p| !p.trim().is_empty())
499 .unwrap_or(&self.default_provider);
500 if let Some(cred) = self.credentials.iter().find(|c| c.provider == provider) {
501 spec.secrets.provider_credentials.push(cred.clone());
502 } else {
503 tracing::warn!(
504 "actor child {}: no credential found for provider '{}'",
505 job.child_session_id,
506 provider
507 );
508 }
509 spec.capabilities.nested_spawn = session.spawn_depth < MAX_SPAWN_DEPTH;
516 spec.capabilities.max_spawn_depth = Some(MAX_SPAWN_DEPTH);
517 spec.capabilities.enforce_permissions = true;
523 spec.capabilities.bypass = session
528 .agent_runtime_state
529 .as_ref()
530 .is_some_and(|s| s.bypass_permissions);
531 spec.capabilities.no_human_approver = session
536 .agent_runtime_state
537 .as_ref()
538 .is_some_and(|s| s.no_human_approver);
539 spec.capabilities.guardian_read_only =
549 session.metadata.get("subagent_type").map(String::as_str) == Some("guardian");
550 if let Some(placement) = self.remote_placements.get(spec.identity.role.as_str()) {
557 spec.placement = Placement::Remote {
558 endpoint: placement.endpoint.clone(),
559 };
560 spec.secrets.worker_auth_token = placement.token.clone();
561 } else if let Some(placement) = self.schedulable_placements.get(spec.identity.role.as_str())
562 {
563 spec.placement = Placement::Schedulable {
571 pool: placement.pool.clone(),
572 };
573 }
574 spec
575 }
576
577 fn placement_stamp_for(&self, spec: &ProvisionSpec) -> Option<String> {
583 let host_label = match &spec.placement {
584 Placement::Remote { .. } => self
585 .remote_placements
586 .get(spec.identity.role.as_str())
587 .and_then(|p| p.host_label.as_deref()),
588 Placement::Schedulable { .. } => self
589 .schedulable_placements
590 .get(spec.identity.role.as_str())
591 .and_then(|p| p.host_label.as_deref()),
592 Placement::Local => None,
593 };
594 placement_metadata(&spec.placement, host_label)
595 }
596
597 async fn resolve_schedulable_worker(
604 &self,
605 role: &str,
606 ) -> std::result::Result<String, AgentError> {
607 let pool = self
608 .schedulable_placements
609 .get(role)
610 .ok_or_else(|| {
611 AgentError::LLM(format!(
612 "schedulable placement for role '{role}' vanished before scheduling"
613 ))
614 })?
615 .pool
616 .clone();
617 let bus = self.bus.as_ref().ok_or_else(|| {
618 AgentError::LLM(format!(
619 "schedulable role '{role}': no mailbox bus configured (subagents.broker)"
620 ))
621 })?;
622
623 let mut q = bamboo_broker::BrokerClient::connect(
626 &bus.endpoint,
627 bamboo_subagent::AgentRef {
628 session_id: format!("sched-q-{role}"),
629 role: None,
630 },
631 &bus.token,
632 )
633 .await
634 .map_err(|e| {
635 AgentError::LLM(format!(
636 "schedulable role '{role}': bus connect failed: {e}"
637 ))
638 })?;
639 let candidates = q.list_connected(&pool).await.map_err(|e| {
640 AgentError::LLM(format!(
641 "schedulable role '{role}': bus presence query failed: {e}"
642 ))
643 })?;
644
645 if candidates.is_empty() {
646 return Err(AgentError::LLM(format!(
647 "schedulable role '{role}': no live worker in pool '{pool}' on the bus \
648 (NOT spawning a local subprocess — a schedulable role has no local fallback)"
649 )));
650 }
651
652 let idx = {
657 let mut cursors = self.schedule_cursor.lock().recover_poison();
658 let cursor = cursors.entry(pool.clone()).or_insert(0);
659 let i = *cursor % candidates.len();
660 *cursor = cursor.wrapping_add(1);
661 i
662 };
663 Ok(candidates[idx].clone())
664 }
665}
666
667#[async_trait]
668impl ExternalChildRunner for ActorChildRunner {
669 async fn should_handle(&self, session: &Session) -> bool {
670 session.metadata.get("runtime.kind") == Some(&"external".to_string())
671 && session.metadata.get("external.protocol") == Some(&"actor".to_string())
672 && session.metadata.get("external.agent_id") == Some(&self.agent_id)
673 }
674
675 fn set_escalation_bridge(&self, bridge: Option<bamboo_subagent::executor::HostBridge>) {
676 *self.escalation_bridge.lock().recover_poison() = bridge;
677 }
678
679 async fn execute_external_child(
680 &self,
681 session: &mut Session,
682 job: &SpawnJob,
683 event_tx: mpsc::Sender<AgentEvent>,
684 cancel_token: CancellationToken,
685 ) -> crate::runtime::runner::Result<()> {
686 let escalation = self.escalation_bridge.lock().recover_poison().clone();
695 let assignment = extract_assignment(session);
696 let mut spec = self.build_spec(session, job);
697 spec.reusable = true;
700 if spec.limits.idle_timeout_secs.is_none() {
701 spec.limits.idle_timeout_secs = Some(POOLED_IDLE_TIMEOUT_SECS);
702 }
703 let pool_key = Self::fingerprint(&spec);
704 let messages: Vec<serde_json::Value> = session
710 .messages
711 .iter()
712 .filter_map(|m| serde_json::to_value(m).ok())
713 .collect();
714 let permission_policy = self.permission_config.as_ref().and_then(|config| {
719 serde_json::to_value(config.to_serializable())
720 .ok()
721 .map(|policy| PermissionPolicyContext {
722 revision: config.policy_revision(),
723 bypass_permissions: session
724 .agent_runtime_state
725 .as_ref()
726 .is_some_and(|state| state.bypass_permissions),
727 session_id: session.id.clone(),
728 workspace_path: session.workspace.clone(),
729 inherit_session_grants: false,
730 policy,
731 })
732 });
733
734 let _slot = self
739 .concurrency
740 .acquire()
741 .await
742 .map_err(|_| AgentError::LLM("actor concurrency limiter closed".to_string()))?;
743
744 let kind = match spec.placement {
752 Placement::Remote { .. } => PlacementKind::Remote,
753 Placement::Schedulable { .. } => PlacementKind::Schedulable,
754 Placement::Local => PlacementKind::Local,
755 };
756 let remote = !matches!(kind, PlacementKind::Local);
757
758 if let Some(placement_meta) = self.placement_stamp_for(&spec) {
764 session
765 .metadata
766 .insert("placement".to_string(), placement_meta);
767 }
768
769 let mut attempt = 0u8;
776 let (result, actor) = loop {
777 let (actor, mut client) = match kind {
778 PlacementKind::Remote => {
779 let placement = self
783 .remote_placements
784 .get(spec.identity.role.as_str())
785 .ok_or_else(|| {
786 AgentError::LLM(format!(
787 "remote placement for role '{}' vanished before connect",
788 spec.identity.role
789 ))
790 })?;
791 let endpoint = placement.endpoint.clone();
792 let trust_cfg = match placement.ca_cert_file.as_deref() {
795 Some(path) => Some(client_config_trusting_cert(path).map_err(|e| {
796 AgentError::LLM(format!(
797 "remote worker CA cert '{}': {e}",
798 path.display()
799 ))
800 })?),
801 None => None,
802 };
803 let client = ChildClient::connect_with_auth_tls(
804 &endpoint,
805 placement.token.as_deref(),
806 trust_cfg,
807 )
808 .await
809 .map_err(|e| {
810 AgentError::LLM(format!("remote actor connect to '{endpoint}' failed: {e}"))
811 })?;
812 let record = AgentRecord {
815 agent_id: job.child_session_id.clone(),
816 role: spec.identity.role.clone(),
817 labels: Vec::new(),
818 endpoint: endpoint.clone(),
819 pid: 0,
820 version: String::new(),
821 started_at: chrono::Utc::now(),
822 lease_expires_at: chrono::Utc::now(),
823 };
824 let _ = endpoint;
825 let actor = PooledWorker {
826 worker: SpawnedChild::remote(record),
827 mailbox_id: job.child_session_id.clone(),
828 };
829 let client: Box<dyn bamboo_subagent::ChildLink> = Box::new(client);
830 (actor, client)
831 }
832 PlacementKind::Schedulable => {
833 let bus = self.bus.as_ref().ok_or_else(|| {
840 AgentError::LLM(
841 "schedulable sub-agents require a mailbox bus (subagents.broker)"
842 .to_string(),
843 )
844 })?;
845 let mailbox_id = self
846 .resolve_schedulable_worker(spec.identity.role.as_str())
847 .await?;
848 let parent = bamboo_subagent::AgentRef {
849 session_id: format!("p-{}", job.child_session_id),
850 role: None,
851 };
852 let link = bamboo_broker::BrokerChildLink::connect(
853 &bus.endpoint,
854 parent,
855 &bus.token,
856 mailbox_id.clone(),
857 )
858 .await
859 .map_err(|e| {
860 AgentError::LLM(format!(
861 "schedulable link connect to '{mailbox_id}' failed: {e}"
862 ))
863 })?;
864 let actor = PooledWorker {
867 worker: SpawnedChild::remote(AgentRecord {
868 agent_id: mailbox_id.clone(),
869 role: spec.identity.role.clone(),
870 labels: Vec::new(),
871 endpoint: bus.endpoint.clone(),
872 pid: 0,
873 version: String::new(),
874 started_at: chrono::Utc::now(),
875 lease_expires_at: chrono::Utc::now(),
876 }),
877 mailbox_id,
878 };
879 let client: Box<dyn bamboo_subagent::ChildLink> = Box::new(link);
880 (actor, client)
881 }
882 PlacementKind::Local => {
883 let bus = self.bus.as_ref().ok_or_else(|| {
890 AgentError::LLM(
891 "local sub-agents require a mailbox bus (subagents.broker); none is \
892 configured and the bus could not be embedded"
893 .to_string(),
894 )
895 })?;
896 let actor = self.acquire_bus_worker(&pool_key, &spec).await?;
897 let parent = bamboo_subagent::AgentRef {
898 session_id: format!("p-{}", job.child_session_id),
899 role: None,
900 };
901 let link = bamboo_broker::BrokerChildLink::connect(
902 &bus.endpoint,
903 parent,
904 &bus.token,
905 actor.mailbox_id.clone(),
906 )
907 .await
908 .map_err(|e| {
909 AgentError::LLM(format!("broker child link connect failed: {e}"))
910 })?;
911 let client: Box<dyn bamboo_subagent::ChildLink> = Box::new(link);
912 (actor, client)
913 }
914 };
915
916 if let Err(e) = client
917 .send(ParentFrame::Run(RunSpec {
918 assignment: assignment.clone(),
920 reasoning_effort: None,
921 permission_policy: permission_policy.clone(),
922 messages: messages.clone(),
923 }))
924 .await
925 {
926 if !remote {
927 actor.worker.kill().await;
928 }
929 return Err(AgentError::LLM(format!("actor run dispatch failed: {e}")));
930 }
931
932 let (live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
936 let live_guard = super::live::register(
937 &job.child_session_id,
938 live_tx,
939 attempt as u32,
940 self.approval_registry.clone(),
941 );
942
943 let result = drive(
944 &mut *client,
945 &job.parent_session_id,
946 &job.child_session_id,
947 attempt as u32,
948 self.approval_registry.as_ref(),
949 self.approval_decider.as_ref(),
950 self.approval_reviewer.as_ref(),
951 escalation.clone(),
952 &event_tx,
953 &cancel_token,
954 &mut live_rx,
955 Some(WORKER_FIRST_FRAME_TIMEOUT),
962 )
963 .await;
964 drop(live_guard);
970 drop(client);
973
974 if attempt == 0 && matches!(result, Err(AgentError::WorkerUnresponsive(_))) {
981 match kind {
982 PlacementKind::Local => {
983 tracing::warn!(
984 "actor child {} got no first frame; reaping the worker and respawning once",
985 job.child_session_id
986 );
987 actor.worker.kill().await;
988 attempt += 1;
989 continue;
990 }
991 PlacementKind::Schedulable => {
992 tracing::warn!(
993 "scheduled actor child {} got no first frame; re-selecting a pool worker",
994 job.child_session_id
995 );
996 drop(actor);
997 attempt += 1;
998 continue;
999 }
1000 PlacementKind::Remote => {}
1001 }
1002 }
1003 break (result, actor);
1004 };
1005
1006 if remote {
1010 drop(actor);
1011 } else {
1012 match &result {
1013 Ok(_) => self.release_bus_worker(&pool_key, actor).await,
1014 Err(_) => actor.worker.kill().await,
1015 }
1016 }
1017
1018 match result {
1022 Ok(Some(text)) => {
1023 if !text.is_empty() {
1024 session.add_message(bamboo_agent_core::Message::assistant(text, None));
1025 }
1026 Ok(())
1027 }
1028 Ok(None) => Ok(()),
1029 Err(e) => Err(e),
1030 }
1031 }
1032}
1033
1034fn placement_metadata(placement: &Placement, host_label: Option<&str>) -> Option<String> {
1040 let value = match placement {
1043 Placement::Local => return None,
1044 Placement::Remote { endpoint } => serde_json::json!({
1045 "kind": "remote",
1046 "host": host_label.map(str::to_string).unwrap_or_else(|| host_of_endpoint(endpoint)),
1047 }),
1048 Placement::Schedulable { pool } => serde_json::json!({
1049 "kind": "remote",
1050 "host": host_label.unwrap_or(pool),
1051 }),
1052 };
1053 serde_json::to_string(&value).ok()
1054}
1055
1056fn host_of_endpoint(endpoint: &str) -> String {
1058 endpoint
1059 .trim()
1060 .trim_start_matches("wss://")
1061 .trim_start_matches("ws://")
1062 .split(['/', ':'])
1063 .next()
1064 .unwrap_or(endpoint)
1065 .to_string()
1066}
1067
1068async fn drive(
1079 client: &mut dyn bamboo_subagent::ChildLink,
1080 parent_session_id: &str,
1081 child_session_id: &str,
1082 child_attempt: u32,
1083 approval_registry: Option<&super::approval_registry::SharedApprovalRegistry>,
1084 approval_decider: Option<&Arc<dyn ChildApprovalDecider>>,
1085 approval_reviewer: Option<&Arc<dyn ChildApprovalReviewer>>,
1086 escalation_bridge: Option<bamboo_subagent::executor::HostBridge>,
1087 event_tx: &mpsc::Sender<AgentEvent>,
1088 cancel_token: &CancellationToken,
1089 live_rx: &mut mpsc::UnboundedReceiver<ParentFrame>,
1090 first_frame_timeout: Option<Duration>,
1091) -> crate::runtime::runner::Result<Option<String>> {
1092 let mut got_first_frame = false;
1099 let mut first_frame_watch = first_frame_timeout.map(|d| Box::pin(tokio::time::sleep(d)));
1100 loop {
1101 tokio::select! {
1102 _ = cancel_token.cancelled() => {
1103 break;
1105 }
1106 _ = async {
1107 match first_frame_watch.as_mut() {
1108 Some(s) => s.as_mut().await,
1109 None => std::future::pending::<()>().await,
1110 }
1111 }, if !got_first_frame => {
1112 return Err(AgentError::WorkerUnresponsive(format!(
1113 "child {child_session_id} produced no frame within {:?}",
1114 first_frame_timeout.unwrap_or_default()
1115 )));
1116 }
1117 Some(frame) = live_rx.recv() => {
1118 if client.send(frame).await.is_err() {
1120 tracing::warn!("live steering frame could not be sent; connection failing");
1121 }
1122 }
1123 frame = client.next_frame() => {
1124 got_first_frame = true;
1127 first_frame_watch = None;
1128 match frame {
1129 Ok(Some(ChildFrame::Event { event })) => {
1130 if let Ok(ev) = serde_json::from_value::<AgentEvent>(event) {
1132 let _ = event_tx.send(ev).await;
1133 }
1134 }
1135 Ok(Some(ChildFrame::ApprovalRequest { id, body })) => {
1136 if let Some(reviewer) = approval_reviewer
1142 .cloned()
1143 .or_else(child_approval_reviewer)
1144 {
1145 let child = child_session_id.to_string();
1153 let parent = parent_session_id.to_string();
1154 let req_id = id.clone();
1155 let body = body.clone();
1156 let registry = approval_registry.cloned();
1157 tokio::spawn(async move {
1158 let approved = tokio::time::timeout(
1159 CHILD_APPROVAL_TIMEOUT,
1160 reviewer.review(&parent, &child, &body),
1161 )
1162 .await
1163 .unwrap_or(false);
1164 super::live::deliver_approval_scoped(
1165 registry.as_ref(),
1166 &child,
1167 child_attempt,
1168 &req_id,
1169 approved,
1170 );
1171 });
1172 } else if approval_decider.is_some() {
1173 let approved =
1177 decide_child_approval(approval_decider, child_session_id, &body)
1178 .await;
1179 if client
1180 .send(ParentFrame::ApprovalReply { id, approved })
1181 .await
1182 .is_err()
1183 {
1184 tracing::warn!(
1185 "failed to answer approval_request; connection failing"
1186 );
1187 }
1188 } else if let Some(host) = escalation_bridge.clone() {
1189 let child = child_session_id.to_string();
1196 let req_id = id.clone();
1197 let body = body.clone();
1198 let registry = approval_registry.cloned();
1199 tokio::spawn(async move {
1200 let approved = match tokio::time::timeout(
1201 CHILD_APPROVAL_TIMEOUT,
1202 host.approval_call(body),
1203 )
1204 .await
1205 {
1206 Ok(Ok(reply)) => reply
1207 .get("approved")
1208 .and_then(|v| v.as_bool())
1209 .unwrap_or(false),
1210 _ => false,
1212 };
1213 super::live::deliver_approval_scoped(
1214 registry.as_ref(),
1215 &child,
1216 child_attempt,
1217 &req_id,
1218 approved,
1219 );
1220 });
1221 } else {
1222 tracing::warn!(
1226 parent_session_id,
1227 child_session_id,
1228 request_id = %id,
1229 "forced-ask request has no parent-agent reviewer; denying"
1230 );
1231 if client
1232 .send(ParentFrame::ApprovalReply {
1233 id,
1234 approved: false,
1235 })
1236 .await
1237 .is_err()
1238 {
1239 tracing::warn!(
1240 "failed to send fail-closed approval reply; connection failing"
1241 );
1242 }
1243 }
1244 }
1245 Ok(Some(ChildFrame::Terminal { status, result, error, .. })) => {
1246 return match status {
1247 TerminalStatus::Completed => Ok(result),
1248 TerminalStatus::Cancelled => Err(AgentError::Cancelled),
1249 TerminalStatus::Error => Err(AgentError::LLM(
1250 error.unwrap_or_else(|| "actor child errored".to_string()),
1251 )),
1252 TerminalStatus::Suspended => Err(AgentError::LLM(
1257 "nested sub-agent suspend received but resume transport is not wired"
1258 .to_string(),
1259 )),
1260 };
1261 }
1262 Ok(None) => {
1263 return Err(AgentError::LLM(
1264 "actor child closed before terminal".to_string(),
1265 ));
1266 }
1267 Err(e) => {
1268 return Err(AgentError::LLM(format!("actor transport error: {e}")));
1269 }
1270 }
1271 }
1272 }
1273 }
1274
1275 let _ = client.send(ParentFrame::Cancel).await;
1277 Err(AgentError::Cancelled)
1278}
1279
1280fn extract_assignment(session: &Session) -> String {
1282 session
1283 .messages
1284 .iter()
1285 .rev()
1286 .find(|m| matches!(m.role, Role::User))
1287 .map(|m| m.content.clone())
1288 .unwrap_or_else(|| {
1289 session
1290 .metadata
1291 .get("title")
1292 .cloned()
1293 .unwrap_or_else(|| "Execute task".to_string())
1294 })
1295}
1296
1297#[cfg(test)]
1298mod tests {
1299 use super::*;
1300
1301 fn spec_with(
1302 role: &str,
1303 provider: &str,
1304 model: &str,
1305 workspace: Option<&str>,
1306 disabled: Option<Vec<&str>>,
1307 ) -> ProvisionSpec {
1308 let mut spec = ProvisionSpec::new(
1309 ChildIdentity {
1310 child_id: "c".into(),
1311 parent_id: None,
1312 project_key: None,
1313 role: role.into(),
1314 depth: 0,
1315 },
1316 ExecutorSpec::Echo,
1317 "/tmp/fab".into(),
1318 );
1319 spec.workspace = workspace.map(|w| w.to_string());
1320 spec.model = Some(ModelRefSpec {
1321 provider: provider.into(),
1322 model: model.into(),
1323 });
1324 spec.disabled_tools = disabled.map(|d| d.into_iter().map(String::from).collect());
1325 spec
1326 }
1327
1328 #[test]
1329 fn fingerprint_matches_interchangeable_children() {
1330 let a = spec_with(
1333 "explorer",
1334 "p",
1335 "m",
1336 Some("/ws"),
1337 Some(vec!["Bash", "Edit"]),
1338 );
1339 let mut b = spec_with(
1340 "explorer",
1341 "p",
1342 "m",
1343 Some("/ws"),
1344 Some(vec!["Edit", "Bash"]),
1345 );
1346 b.identity.child_id = "other".into();
1347 assert_eq!(
1348 ActorChildRunner::fingerprint(&a),
1349 ActorChildRunner::fingerprint(&b)
1350 );
1351 }
1352
1353 #[test]
1354 fn fingerprint_separates_distinct_runtimes() {
1355 let base = spec_with("explorer", "p", "m", Some("/ws"), None);
1356 let base_fp = ActorChildRunner::fingerprint(&base);
1357 assert_ne!(
1359 base_fp,
1360 ActorChildRunner::fingerprint(&spec_with("writer", "p", "m", Some("/ws"), None))
1361 );
1362 assert_ne!(
1363 base_fp,
1364 ActorChildRunner::fingerprint(&spec_with("explorer", "p2", "m", Some("/ws"), None))
1365 );
1366 assert_ne!(
1367 base_fp,
1368 ActorChildRunner::fingerprint(&spec_with("explorer", "p", "m2", Some("/ws"), None))
1369 );
1370 assert_ne!(
1371 base_fp,
1372 ActorChildRunner::fingerprint(&spec_with("explorer", "p", "m", Some("/ws2"), None))
1373 );
1374 assert_ne!(
1375 base_fp,
1376 ActorChildRunner::fingerprint(&spec_with(
1377 "explorer",
1378 "p",
1379 "m",
1380 Some("/ws"),
1381 Some(vec!["Bash"])
1382 ))
1383 );
1384 }
1385
1386 #[test]
1387 fn fingerprint_splits_on_baked_capabilities() {
1388 let base_fp =
1393 ActorChildRunner::fingerprint(&spec_with("explorer", "p", "m", Some("/ws"), None));
1394
1395 let mut depth = spec_with("explorer", "p", "m", Some("/ws"), None);
1396 depth.identity.depth = 2;
1397 assert_ne!(
1398 base_fp,
1399 ActorChildRunner::fingerprint(&depth),
1400 "depth must split"
1401 );
1402
1403 let mut nested = spec_with("explorer", "p", "m", Some("/ws"), None);
1404 nested.capabilities.nested_spawn = true;
1405 assert_ne!(
1406 base_fp,
1407 ActorChildRunner::fingerprint(&nested),
1408 "nested_spawn must split"
1409 );
1410
1411 let mut bypass = spec_with("explorer", "p", "m", Some("/ws"), None);
1412 bypass.capabilities.bypass = true;
1413 assert_ne!(
1414 base_fp,
1415 ActorChildRunner::fingerprint(&bypass),
1416 "bypass must split"
1417 );
1418
1419 let mut enforce = spec_with("explorer", "p", "m", Some("/ws"), None);
1420 enforce.capabilities.enforce_permissions = true;
1421 assert_ne!(
1422 base_fp,
1423 ActorChildRunner::fingerprint(&enforce),
1424 "enforce_permissions must split"
1425 );
1426
1427 let mut cap = spec_with("explorer", "p", "m", Some("/ws"), None);
1428 cap.capabilities.max_spawn_depth = Some(8);
1429 assert_ne!(
1430 base_fp,
1431 ActorChildRunner::fingerprint(&cap),
1432 "max_spawn_depth must split"
1433 );
1434
1435 let mut nha = spec_with("explorer", "p", "m", Some("/ws"), None);
1439 nha.capabilities.no_human_approver = true;
1440 assert_ne!(
1441 base_fp,
1442 ActorChildRunner::fingerprint(&nha),
1443 "no_human_approver must split"
1444 );
1445
1446 let mut gro = spec_with("explorer", "p", "m", Some("/ws"), None);
1449 gro.capabilities.guardian_read_only = true;
1450 assert_ne!(
1451 base_fp,
1452 ActorChildRunner::fingerprint(&gro),
1453 "guardian_read_only must split"
1454 );
1455 }
1456
1457 struct StaticDecider(bool);
1458
1459 #[async_trait]
1460 impl ChildApprovalDecider for StaticDecider {
1461 async fn decide(&self, _child: &str, _req: &serde_json::Value) -> bool {
1462 self.0
1463 }
1464 }
1465
1466 struct RecordingReviewer {
1467 reviewed: mpsc::UnboundedSender<(String, String, serde_json::Value)>,
1468 }
1469
1470 #[async_trait]
1471 impl ChildApprovalReviewer for RecordingReviewer {
1472 async fn review(&self, parent: &str, child: &str, request: &serde_json::Value) -> bool {
1473 let _ = self
1474 .reviewed
1475 .send((parent.to_string(), child.to_string(), request.clone()));
1476 true
1477 }
1478 }
1479
1480 struct SilentLink;
1485 #[async_trait]
1486 impl bamboo_subagent::ChildLink for SilentLink {
1487 async fn send(&mut self, _: ParentFrame) -> bamboo_subagent::TransportResult<()> {
1488 Ok(())
1489 }
1490 async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
1491 std::future::pending().await
1492 }
1493 }
1494
1495 struct InstantTerminalLink {
1497 done: bool,
1498 }
1499
1500 struct ApprovalRoundTripLink {
1501 step: u8,
1502 approval_reply: Option<(String, bool)>,
1503 }
1504
1505 #[async_trait]
1506 impl bamboo_subagent::ChildLink for ApprovalRoundTripLink {
1507 async fn send(&mut self, frame: ParentFrame) -> bamboo_subagent::TransportResult<()> {
1508 if let ParentFrame::ApprovalReply { id, approved } = frame {
1509 self.approval_reply = Some((id, approved));
1510 self.step = 2;
1511 }
1512 Ok(())
1513 }
1514
1515 async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
1516 match self.step {
1517 0 => {
1518 self.step = 1;
1519 Ok(Some(ChildFrame::ApprovalRequest {
1520 id: "approval-1".into(),
1521 body: serde_json::json!({
1522 "tool_name": "Bash",
1523 "permission": "execute",
1524 "resource": "rm -rf target",
1525 "permission_request": {"reason_code": "hard_dangerous"}
1526 }),
1527 }))
1528 }
1529 1 => std::future::pending().await,
1530 2 => {
1531 self.step = 3;
1532 Ok(Some(ChildFrame::Terminal {
1533 status: TerminalStatus::Completed,
1534 result: Some("done".into()),
1535 error: None,
1536 transcript: vec![],
1537 }))
1538 }
1539 _ => std::future::pending().await,
1540 }
1541 }
1542 }
1543
1544 #[tokio::test]
1545 async fn drive_routes_forced_ask_to_parent_reviewer_without_human_event() {
1546 let (event_tx, mut event_rx) = mpsc::channel::<AgentEvent>(8);
1547 let (review_tx, mut review_rx) = mpsc::unbounded_channel();
1548 let reviewer: Arc<dyn ChildApprovalReviewer> = Arc::new(RecordingReviewer {
1549 reviewed: review_tx,
1550 });
1551 let cancel = CancellationToken::new();
1552 let (live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
1553 let live_guard = crate::external_agents::live::register("child-reviewer", live_tx, 0, None);
1554 let mut link = ApprovalRoundTripLink {
1555 step: 0,
1556 approval_reply: None,
1557 };
1558
1559 let result = tokio::time::timeout(
1560 Duration::from_secs(1),
1561 drive(
1562 &mut link,
1563 "parent-reviewer",
1564 "child-reviewer",
1565 0,
1566 None,
1567 None,
1568 Some(&reviewer),
1569 None,
1570 &event_tx,
1571 &cancel,
1572 &mut live_rx,
1573 None,
1574 ),
1575 )
1576 .await
1577 .expect("worker must receive the reviewer verdict before terminating");
1578
1579 assert_eq!(result.ok().flatten().as_deref(), Some("done"));
1580 assert_eq!(
1581 link.approval_reply,
1582 Some(("approval-1".to_string(), true)),
1583 "reviewer verdict must traverse the live route back to the worker"
1584 );
1585 let (parent, child, body) = tokio::time::timeout(Duration::from_secs(1), review_rx.recv())
1586 .await
1587 .expect("reviewer should be invoked off-loop")
1588 .expect("review channel should remain open");
1589 assert_eq!(parent, "parent-reviewer");
1590 assert_eq!(child, "child-reviewer");
1591 assert_eq!(
1592 body.pointer("/permission_request/reason_code")
1593 .and_then(serde_json::Value::as_str),
1594 Some("hard_dangerous")
1595 );
1596 assert!(
1597 event_rx.try_recv().is_err(),
1598 "must not emit a human-review event"
1599 );
1600 drop(live_guard);
1601 }
1602
1603 #[tokio::test]
1604 async fn drive_denies_forced_ask_without_parent_reviewer_or_manual_event() {
1605 let (event_tx, mut event_rx) = mpsc::channel::<AgentEvent>(8);
1606 let cancel = CancellationToken::new();
1607 let (_live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
1608 let mut link = ApprovalRoundTripLink {
1609 step: 0,
1610 approval_reply: None,
1611 };
1612
1613 let result = tokio::time::timeout(
1614 Duration::from_secs(1),
1615 drive(
1616 &mut link,
1617 "parent-no-reviewer",
1618 "child-no-reviewer",
1619 0,
1620 None,
1621 None,
1622 None,
1623 None,
1624 &event_tx,
1625 &cancel,
1626 &mut live_rx,
1627 None,
1628 ),
1629 )
1630 .await
1631 .expect("fail-closed reply must unblock the child immediately");
1632
1633 assert_eq!(result.ok().flatten().as_deref(), Some("done"));
1634 assert_eq!(link.approval_reply, Some(("approval-1".to_string(), false)));
1635 assert!(
1636 event_rx.try_recv().is_err(),
1637 "missing parent review must not surface a manual approval event"
1638 );
1639 }
1640 #[async_trait]
1641 impl bamboo_subagent::ChildLink for InstantTerminalLink {
1642 async fn send(&mut self, _: ParentFrame) -> bamboo_subagent::TransportResult<()> {
1643 Ok(())
1644 }
1645 async fn next_frame(&mut self) -> bamboo_subagent::TransportResult<Option<ChildFrame>> {
1646 if self.done {
1647 std::future::pending().await
1648 } else {
1649 self.done = true;
1650 Ok(Some(ChildFrame::Terminal {
1651 status: TerminalStatus::Completed,
1652 result: Some("done".into()),
1653 error: None,
1654 transcript: vec![],
1655 }))
1656 }
1657 }
1658 }
1659
1660 #[tokio::test]
1661 async fn drive_trips_first_frame_watchdog_on_a_silent_worker() {
1662 let (event_tx, _rx) = mpsc::channel::<AgentEvent>(8);
1663 let cancel = CancellationToken::new();
1664 let (_live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
1665 let mut link = SilentLink;
1666 let r = drive(
1667 &mut link,
1668 "parent-x",
1669 "child-x",
1670 0,
1671 None,
1672 None,
1673 None,
1674 None,
1675 &event_tx,
1676 &cancel,
1677 &mut live_rx,
1678 Some(Duration::from_millis(100)),
1679 )
1680 .await;
1681 assert!(
1682 matches!(r, Err(AgentError::WorkerUnresponsive(_))),
1683 "a silent worker must trip the first-frame watchdog, got {r:?}"
1684 );
1685 }
1686
1687 #[tokio::test]
1688 async fn drive_does_not_trip_when_a_frame_arrives() {
1689 let (event_tx, _rx) = mpsc::channel::<AgentEvent>(8);
1690 let cancel = CancellationToken::new();
1691 let (_live_tx, mut live_rx) = mpsc::unbounded_channel::<ParentFrame>();
1692 let mut link = InstantTerminalLink { done: false };
1693 let r = drive(
1696 &mut link,
1697 "parent-y",
1698 "child-y",
1699 0,
1700 None,
1701 None,
1702 None,
1703 None,
1704 &event_tx,
1705 &cancel,
1706 &mut live_rx,
1707 Some(Duration::from_millis(50)),
1708 )
1709 .await;
1710 assert_eq!(r.ok().flatten().as_deref(), Some("done"));
1711 }
1712
1713 #[tokio::test]
1714 async fn child_approval_fails_closed_without_decider() {
1715 let body = serde_json::json!({"tool_name":"Bash","permission":"run","resource":"rm -rf /"});
1717 assert!(!decide_child_approval(None, "child-1", &body).await);
1718 }
1719
1720 #[tokio::test]
1721 async fn child_approval_honors_wired_decider() {
1722 let body =
1723 serde_json::json!({"tool_name":"Write","permission":"write","resource":"/tmp/x"});
1724 let approve: Arc<dyn ChildApprovalDecider> = Arc::new(StaticDecider(true));
1725 let deny: Arc<dyn ChildApprovalDecider> = Arc::new(StaticDecider(false));
1726 assert!(decide_child_approval(Some(&approve), "child-1", &body).await);
1727 assert!(!decide_child_approval(Some(&deny), "child-1", &body).await);
1728 }
1729
1730 use crate::runtime::execution::SpawnJob;
1733 use bamboo_agent_core::Session;
1734
1735 fn bogus_runner(placements: HashMap<String, ResolvedRemotePlacement>) -> ActorChildRunner {
1738 ActorChildRunner::new(
1739 "test-actor".into(),
1740 PathBuf::from("/bin/false"),
1741 vec![],
1742 std::env::temp_dir().join("bamboo-test-fab-193"),
1743 ExecutorSpec::Echo,
1744 vec![],
1745 "anthropic".into(),
1746 4,
1747 )
1748 .with_remote_placements(placements)
1749 }
1750
1751 fn session_of_role(role: &str, assignment: &str) -> Session {
1754 let mut s = Session::new("child-1", "test-model");
1755 s.metadata
1756 .insert("subagent_type".to_string(), role.to_string());
1757 s.add_message(bamboo_agent_core::Message::user(assignment));
1758 s
1759 }
1760
1761 fn job_for(child: &str) -> SpawnJob {
1762 SpawnJob {
1763 parent_session_id: "parent-1".into(),
1764 child_session_id: child.into(),
1765 model: String::new(),
1766 disabled_tools: None,
1767 }
1768 }
1769
1770 #[derive(Default)]
1771 struct RecordingChildSessionPort {
1772 saved: std::sync::Mutex<Option<Session>>,
1773 }
1774
1775 impl RecordingChildSessionPort {
1776 fn saved_child(&self) -> Session {
1777 self.saved
1778 .lock()
1779 .expect("saved-child fixture lock")
1780 .clone()
1781 .expect("create_child_action must save the child")
1782 }
1783 }
1784
1785 #[async_trait]
1786 impl crate::session_app::child_session::ChildSessionPort for RecordingChildSessionPort {
1787 async fn load_root_session(
1788 &self,
1789 _root_id: &str,
1790 ) -> Result<Session, crate::session_app::child_session::ChildSessionError> {
1791 unreachable!("create_child_action does not load the root")
1792 }
1793
1794 async fn load_child_for_parent(
1795 &self,
1796 _parent_id: &str,
1797 _child_id: &str,
1798 ) -> Result<Session, crate::session_app::child_session::ChildSessionError> {
1799 unreachable!("create_child_action does not reload the child")
1800 }
1801
1802 async fn save_child_session(
1803 &self,
1804 child: &mut Session,
1805 ) -> Result<(), crate::session_app::child_session::ChildSessionError> {
1806 *self.saved.lock().expect("saved-child fixture lock") = Some(child.clone());
1807 Ok(())
1808 }
1809
1810 async fn save_child_session_authoritative_flags(
1811 &self,
1812 _child: &mut Session,
1813 ) -> Result<(), crate::session_app::child_session::ChildSessionError> {
1814 unreachable!("new-child creation uses the ordinary save")
1815 }
1816
1817 async fn is_child_running(&self, _child_id: &str) -> bool {
1818 false
1819 }
1820
1821 async fn list_children(
1822 &self,
1823 _parent_id: &str,
1824 ) -> Vec<crate::session_app::child_session::ChildSessionEntry> {
1825 Vec::new()
1826 }
1827
1828 async fn enqueue_child_run(
1829 &self,
1830 _parent: &Session,
1831 _child: &Session,
1832 ) -> Result<(), crate::session_app::child_session::ChildSessionError> {
1833 unreachable!("fixture creates the child with auto_run=false")
1834 }
1835
1836 async fn cancel_child_run_and_wait(
1837 &self,
1838 _child_id: &str,
1839 ) -> Result<(), crate::session_app::child_session::ChildSessionError> {
1840 unreachable!("create_child_action does not cancel")
1841 }
1842
1843 async fn delete_child_session(
1844 &self,
1845 _parent_id: &str,
1846 _child_id: &str,
1847 ) -> Result<
1848 crate::session_app::child_session::DeleteChildResult,
1849 crate::session_app::child_session::ChildSessionError,
1850 > {
1851 unreachable!("create_child_action does not delete")
1852 }
1853
1854 async fn get_child_runner_info(
1855 &self,
1856 _child_id: &str,
1857 ) -> Option<crate::session_app::child_session::ChildRunnerInfo> {
1858 None
1859 }
1860
1861 async fn register_parent_wait_for_child(
1862 &self,
1863 _parent_session_id: &str,
1864 _child_session_id: &str,
1865 _tool_call_id: Option<&str>,
1866 ) -> Result<(), crate::session_app::child_session::ChildSessionError> {
1867 unreachable!("create_child_action does not register a wait")
1868 }
1869
1870 async fn register_parent_wait_for_children(
1871 &self,
1872 _parent_session_id: &str,
1873 _child_session_ids: &[String],
1874 _policy: bamboo_domain::session::runtime_state::ChildWaitPolicy,
1875 ) -> Result<usize, crate::session_app::child_session::ChildSessionError> {
1876 unreachable!("create_child_action does not register a wait")
1877 }
1878
1879 async fn active_child_ids(&self, _parent_session_id: &str) -> Vec<String> {
1880 Vec::new()
1881 }
1882
1883 async fn find_resident_child(
1884 &self,
1885 _root_session_id: &str,
1886 _resident_name: &str,
1887 ) -> Option<String> {
1888 None
1889 }
1890
1891 async fn ensure_child_indexed(&self, _child_session_id: &str) {}
1892 }
1893
1894 #[test]
1895 fn build_spec_sets_remote_placement_for_matching_role() {
1896 let mut placements = HashMap::new();
1897 placements.insert(
1898 "explorer".to_string(),
1899 ResolvedRemotePlacement {
1900 endpoint: "wss://gpu-host:8443".into(),
1901 token: Some("T-secret".into()),
1902 ca_cert_file: None,
1903 host_label: None,
1904 },
1905 );
1906 let runner = bogus_runner(placements);
1907
1908 let s = session_of_role("explorer", "do the thing");
1910 let spec = runner.build_spec(&s, &job_for("child-1"));
1911 match &spec.placement {
1912 Placement::Remote { endpoint } => assert_eq!(endpoint, "wss://gpu-host:8443"),
1913 other => panic!("expected Remote, got {other:?}"),
1914 }
1915 assert_eq!(spec.secrets.worker_auth_token.as_deref(), Some("T-secret"));
1916 }
1917
1918 #[test]
1919 fn build_spec_leaves_local_for_unmatched_role() {
1920 let mut placements = HashMap::new();
1921 placements.insert(
1922 "explorer".to_string(),
1923 ResolvedRemotePlacement {
1924 endpoint: "wss://gpu-host:8443".into(),
1925 token: Some("T".into()),
1926 ca_cert_file: None,
1927 host_label: None,
1928 },
1929 );
1930 let runner = bogus_runner(placements);
1931
1932 let s = session_of_role("writer", "do the thing");
1934 let spec = runner.build_spec(&s, &job_for("child-1"));
1935 assert_eq!(spec.placement, Placement::Local);
1936 assert!(spec.secrets.worker_auth_token.is_none());
1937 }
1938
1939 #[test]
1940 fn build_spec_local_when_no_placements() {
1941 let runner = bogus_runner(HashMap::new());
1942 let s = session_of_role("explorer", "do the thing");
1943 let spec = runner.build_spec(&s, &job_for("child-1"));
1944 assert_eq!(spec.placement, Placement::Local);
1945 assert!(spec.secrets.worker_auth_token.is_none());
1946 }
1947
1948 #[tokio::test]
1949 async fn build_spec_preserves_inherited_bypass_for_child_worker() {
1950 let runner = bogus_runner(HashMap::new());
1954 let mut parent = Session::new("parent-bypass", "test-model");
1955 parent
1956 .agent_runtime_state
1957 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
1958 .bypass_permissions = true;
1959 let workspace = tempfile::tempdir().expect("workspace fixture");
1960 let port = RecordingChildSessionPort::default();
1961 let child_id = format!("child-bypass-{}", uuid::Uuid::new_v4());
1962 crate::session_app::child_session::create_child_action(
1963 &port,
1964 crate::session_app::child_session::CreateChildInput {
1965 parent_session: parent,
1966 child_id: child_id.clone(),
1967 title: "Bypassed child".to_string(),
1968 responsibility: "Run ordinary commands".to_string(),
1969 assignment_prompt: "run an ordinary command".to_string(),
1970 subagent_type: "explorer".to_string(),
1971 workspace: workspace.path().to_string_lossy().into_owned(),
1972 model_override: None,
1973 model_ref_override: None,
1974 runtime_metadata: HashMap::new(),
1975 auto_run: false,
1976 reasoning_effort: None,
1977 lifecycle: None,
1978 resident_name: None,
1979 resident_context: None,
1980 disabled_tools: None,
1981 context_fork: None,
1982 },
1983 )
1984 .await
1985 .expect("create inherited-bypass child");
1986 let child = port.saved_child();
1987
1988 assert!(
1989 child
1990 .agent_runtime_state
1991 .as_ref()
1992 .is_some_and(|state| state.bypass_permissions),
1993 "create_child_action must inherit bypass from the parent"
1994 );
1995
1996 let spec = runner.build_spec(&child, &job_for(&child_id));
1997
1998 assert!(spec.capabilities.bypass, "child worker must inherit bypass");
1999 assert!(
2000 spec.capabilities.enforce_permissions,
2001 "forced-ask evaluation must remain active under bypass"
2002 );
2003 }
2004
2005 #[test]
2006 fn placement_metadata_stamps_remote_and_schedulable_not_local() {
2007 assert_eq!(placement_metadata(&Placement::Local, None), None);
2009
2010 let r = placement_metadata(
2012 &Placement::Remote {
2013 endpoint: "wss://10.0.0.5:8443/stream".into(),
2014 },
2015 None,
2016 )
2017 .unwrap();
2018 assert!(r.contains(r#""kind":"remote""#), "{r}");
2019 assert!(r.contains(r#""host":"10.0.0.5""#), "{r}");
2020
2021 let labeled = placement_metadata(
2023 &Placement::Remote {
2024 endpoint: "ws://169.254.230.101:8899".into(),
2025 },
2026 Some("mini"),
2027 )
2028 .unwrap();
2029 assert!(labeled.contains(r#""host":"mini""#), "{labeled}");
2030
2031 let s = placement_metadata(
2033 &Placement::Schedulable {
2034 pool: "explorers".into(),
2035 },
2036 Some("mini"),
2037 )
2038 .unwrap();
2039 assert!(s.contains(r#""kind":"remote""#), "{s}");
2040 assert!(s.contains(r#""host":"mini""#), "{s}");
2041
2042 let p: bamboo_storage::SessionPlacement = serde_json::from_str(&labeled).unwrap();
2044 assert_eq!(p.kind, "remote");
2045 assert_eq!(p.host, "mini");
2046 }
2047
2048 #[tokio::test]
2055 async fn execute_external_child_routes_role_to_remote_worker_without_spawning() {
2056 let token = "remote-test-token";
2058 let server = bamboo_subagent::transport::WsServer::bind_with_token(
2059 (std::net::Ipv4Addr::LOCALHOST, 0).into(),
2060 Some(token.to_string()),
2061 )
2062 .await
2063 .expect("bind resident worker");
2064 let endpoint = server.ws_endpoint(); let srv = tokio::spawn(async move {
2066 let _ = server
2068 .serve(Arc::new(bamboo_subagent::executor::EchoExecutor))
2069 .await;
2070 });
2071
2072 let mut placements = HashMap::new();
2074 placements.insert(
2075 "explorer".to_string(),
2076 ResolvedRemotePlacement {
2077 endpoint: endpoint.clone(),
2078 token: Some(token.to_string()),
2079 ca_cert_file: None,
2080 host_label: Some("mini-e2e".into()), },
2082 );
2083 let runner = bogus_runner(placements);
2084
2085 let mut session = session_of_role("explorer", "hello remote");
2087 let job = job_for("child-1");
2088 let (event_tx, mut event_rx) = mpsc::channel::<AgentEvent>(64);
2089 let cancel = CancellationToken::new();
2090
2091 let result = tokio::time::timeout(
2092 Duration::from_secs(10),
2093 runner.execute_external_child(&mut session, &job, event_tx, cancel),
2094 )
2095 .await
2096 .expect("run did not hang")
2097 .expect("remote run succeeded (connected to resident worker, did not spawn)");
2098
2099 let _ = result;
2100 let last = session
2103 .messages
2104 .iter()
2105 .rev()
2106 .find(|m| matches!(m.role, Role::Assistant))
2107 .expect("an assistant reply was written back");
2108 assert!(
2109 last.content.contains("echo:"),
2110 "expected echo reply, got {:?}",
2111 last.content
2112 );
2113
2114 let placement = session
2117 .metadata
2118 .get("placement")
2119 .expect("remote child session stamped with a placement");
2120 assert!(placement.contains(r#""kind":"remote""#), "{placement}");
2121 assert!(placement.contains(r#""host":"mini-e2e""#), "{placement}");
2122
2123 let mut saw_event = false;
2126 while let Ok(Some(_ev)) =
2127 tokio::time::timeout(Duration::from_millis(50), event_rx.recv()).await
2128 {
2129 saw_event = true;
2130 }
2131 let _ = saw_event;
2132
2133 srv.abort();
2134 }
2135
2136 fn bogus_sched_runner(
2142 remote: HashMap<String, ResolvedRemotePlacement>,
2143 sched: HashMap<String, ResolvedSchedulablePlacement>,
2144 ) -> ActorChildRunner {
2145 ActorChildRunner::new(
2146 "test-actor".into(),
2147 PathBuf::from("/bin/false"),
2148 vec![],
2149 std::env::temp_dir().join("bamboo-test-fab-181"),
2150 ExecutorSpec::Echo,
2151 vec![],
2152 "anthropic".into(),
2153 4,
2154 )
2155 .with_remote_placements(remote)
2156 .with_schedulable_placements(sched)
2157 }
2158
2159 fn sched_placement(
2160 pool: &str,
2161 _registry_url: impl Into<String>,
2162 ) -> ResolvedSchedulablePlacement {
2163 ResolvedSchedulablePlacement {
2164 pool: pool.into(),
2165 host_label: None,
2166 }
2167 }
2168
2169 #[test]
2170 fn build_spec_sets_schedulable_placement_for_matching_role() {
2171 let mut sched = HashMap::new();
2172 sched.insert(
2173 "explorer".to_string(),
2174 sched_placement("gpu-pool", "unused"),
2175 );
2176 let runner = bogus_sched_runner(HashMap::new(), sched);
2177
2178 let s = session_of_role("explorer", "do the thing");
2179 let spec = runner.build_spec(&s, &job_for("child-1"));
2180 match &spec.placement {
2181 Placement::Schedulable { pool } => assert_eq!(pool, "gpu-pool"),
2182 other => panic!("expected Schedulable, got {other:?}"),
2183 }
2184 assert!(spec.secrets.worker_auth_token.is_none());
2186 }
2187
2188 #[test]
2189 fn build_spec_remote_wins_when_role_in_both_maps() {
2190 let mut remote = HashMap::new();
2193 remote.insert(
2194 "explorer".to_string(),
2195 ResolvedRemotePlacement {
2196 endpoint: "wss://fixed-host:8443".into(),
2197 token: Some("T-remote".into()),
2198 ca_cert_file: None,
2199 host_label: None,
2200 },
2201 );
2202 let mut sched = HashMap::new();
2203 sched.insert(
2204 "explorer".to_string(),
2205 sched_placement("gpu-pool", "https://control-plane:9562"),
2206 );
2207 let runner = bogus_sched_runner(remote, sched);
2208
2209 let s = session_of_role("explorer", "do the thing");
2210 let spec = runner.build_spec(&s, &job_for("child-1"));
2211 match &spec.placement {
2212 Placement::Remote { endpoint } => assert_eq!(endpoint, "wss://fixed-host:8443"),
2213 other => panic!("expected Remote (precedence), got {other:?}"),
2214 }
2215 assert_eq!(spec.secrets.worker_auth_token.as_deref(), Some("T-remote"));
2216 }
2217
2218 #[test]
2219 fn build_spec_local_for_unmatched_schedulable_role() {
2220 let mut sched = HashMap::new();
2221 sched.insert(
2222 "explorer".to_string(),
2223 sched_placement("gpu-pool", "https://control-plane:9562"),
2224 );
2225 let runner = bogus_sched_runner(HashMap::new(), sched);
2226 let s = session_of_role("writer", "do the thing");
2227 let spec = runner.build_spec(&s, &job_for("child-1"));
2228 assert_eq!(spec.placement, Placement::Local);
2229 assert!(spec.secrets.worker_auth_token.is_none());
2230 }
2231
2232 #[test]
2237 fn placement_stamp_uses_node_label_for_remote_and_schedulable() {
2238 let mut remote = HashMap::new();
2240 remote.insert(
2241 "explorer".to_string(),
2242 ResolvedRemotePlacement {
2243 endpoint: "ws://169.254.230.101:8899".into(),
2244 token: None,
2245 ca_cert_file: None,
2246 host_label: Some("mini".into()),
2247 },
2248 );
2249 let runner = bogus_runner(remote);
2250 let spec = runner.build_spec(&session_of_role("explorer", "go"), &job_for("c1"));
2251 let stamp = runner
2252 .placement_stamp_for(&spec)
2253 .expect("remote child is stamped");
2254 assert!(stamp.contains(r#""kind":"remote""#), "{stamp}");
2255 assert!(stamp.contains(r#""host":"mini""#), "{stamp}");
2256
2257 let mut remote_nolabel = HashMap::new();
2259 remote_nolabel.insert(
2260 "explorer".to_string(),
2261 ResolvedRemotePlacement {
2262 endpoint: "ws://169.254.230.101:8899".into(),
2263 token: None,
2264 ca_cert_file: None,
2265 host_label: None,
2266 },
2267 );
2268 let r2 = bogus_runner(remote_nolabel);
2269 let spec2 = r2.build_spec(&session_of_role("explorer", "go"), &job_for("c1"));
2270 assert!(r2
2271 .placement_stamp_for(&spec2)
2272 .unwrap()
2273 .contains(r#""host":"169.254.230.101""#));
2274
2275 let mut sched = HashMap::new();
2277 sched.insert(
2278 "mac-mini-monitor".to_string(),
2279 ResolvedSchedulablePlacement {
2280 pool: "mac-mini-monitor".into(),
2281 host_label: Some("mini".into()),
2282 },
2283 );
2284 let sr = bogus_sched_runner(HashMap::new(), sched);
2285 let spec3 = sr.build_spec(&session_of_role("mac-mini-monitor", "go"), &job_for("c1"));
2286 let stamp3 = sr
2287 .placement_stamp_for(&spec3)
2288 .expect("scheduled child is stamped");
2289 assert!(stamp3.contains(r#""kind":"remote""#), "{stamp3}");
2290 assert!(stamp3.contains(r#""host":"mini""#), "{stamp3}");
2291
2292 let local = bogus_runner(HashMap::new());
2294 let spec4 = local.build_spec(&session_of_role("writer", "go"), &job_for("c1"));
2295 assert_eq!(local.placement_stamp_for(&spec4), None);
2296 }
2297
2298 async fn start_bus() -> (String, tempfile::TempDir) {
2301 let dir = tempfile::tempdir().unwrap();
2302 let core = std::sync::Arc::new(bamboo_broker::BrokerCore::new(dir.path()));
2303 let server = std::sync::Arc::new(bamboo_broker::BrokerServer::new(core, "t"));
2304 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2305 let addr = listener.local_addr().unwrap();
2306 tokio::spawn(async move {
2307 let _ = server.serve(listener).await;
2308 });
2309 (format!("ws://{addr}"), dir)
2310 }
2311
2312 async fn join_pool(endpoint: &str, id: &str, pool: &str) -> bamboo_broker::BrokerClient {
2313 let mut c = bamboo_broker::BrokerClient::connect(
2314 endpoint,
2315 bamboo_subagent::AgentRef {
2316 session_id: id.into(),
2317 role: Some(pool.into()),
2318 },
2319 "t",
2320 )
2321 .await
2322 .unwrap();
2323 c.subscribe().await.unwrap();
2324 c
2325 }
2326
2327 fn sched_runner_on_bus(endpoint: &str, child_role: &str, pool: &str) -> ActorChildRunner {
2328 let mut sched = HashMap::new();
2329 sched.insert(child_role.to_string(), sched_placement(pool, "unused"));
2330 bogus_sched_runner(HashMap::new(), sched).with_bus(Some(bamboo_subagent::BusEndpoint {
2331 endpoint: endpoint.into(),
2332 token: "t".into(),
2333 }))
2334 }
2335
2336 #[tokio::test]
2337 async fn resolve_schedulable_picks_a_live_bus_worker() {
2338 let (endpoint, _dir) = start_bus().await;
2339 let _w = join_pool(&endpoint, "w-gpu", "gpu-pool").await;
2340 let runner = sched_runner_on_bus(&endpoint, "explorer", "gpu-pool");
2341
2342 let mailbox = runner
2343 .resolve_schedulable_worker("explorer")
2344 .await
2345 .expect("a live pool worker is found on the bus");
2346 assert_eq!(mailbox, "w-gpu");
2347 }
2348
2349 #[tokio::test]
2350 async fn resolve_schedulable_round_robins_over_pool_workers() {
2351 let (endpoint, _dir) = start_bus().await;
2352 let _a = join_pool(&endpoint, "w-a", "gpu-pool").await;
2353 let _b = join_pool(&endpoint, "w-b", "gpu-pool").await;
2354 let runner = sched_runner_on_bus(&endpoint, "explorer", "gpu-pool");
2355
2356 let mut picked = std::collections::HashSet::new();
2358 for _ in 0..6 {
2359 picked.insert(runner.resolve_schedulable_worker("explorer").await.unwrap());
2360 }
2361 assert_eq!(
2362 picked,
2363 ["w-a".to_string(), "w-b".to_string()].into_iter().collect(),
2364 "round-robin must cover every connected pool worker"
2365 );
2366 }
2367
2368 #[tokio::test]
2369 async fn resolve_schedulable_errors_on_empty_pool() {
2370 let (endpoint, _dir) = start_bus().await;
2371 let runner = sched_runner_on_bus(&endpoint, "explorer", "gpu-pool");
2373
2374 let err = runner
2375 .resolve_schedulable_worker("explorer")
2376 .await
2377 .expect_err("an empty pool is terminal — no local fallback")
2378 .to_string();
2379 assert!(err.contains("no live worker in pool"), "got: {err}");
2380 assert!(err.contains("NOT spawning"), "got: {err}");
2381 }
2382
2383 #[tokio::test]
2390 async fn execute_external_child_runs_schedulable_over_bus_and_stamps_node_label() {
2391 let (endpoint, _dir) = start_bus().await;
2392
2393 let ep = endpoint.clone();
2395 let worker = tokio::spawn(async move {
2396 let _ = bamboo_broker::serve_executor(
2397 &ep,
2398 bamboo_subagent::AgentRef {
2399 session_id: "mmm-worker".into(),
2400 role: Some("mac-mini-monitor".into()),
2401 },
2402 "t",
2403 std::sync::Arc::new(bamboo_subagent::executor::EchoExecutor),
2404 )
2405 .await;
2406 });
2407
2408 let mut probe = bamboo_broker::BrokerClient::connect(
2411 &endpoint,
2412 bamboo_subagent::AgentRef {
2413 session_id: "probe".into(),
2414 role: None,
2415 },
2416 "t",
2417 )
2418 .await
2419 .unwrap();
2420 let mut ready = false;
2421 for _ in 0..100 {
2422 if probe
2423 .list_connected("mac-mini-monitor")
2424 .await
2425 .unwrap()
2426 .iter()
2427 .any(|id| id == "mmm-worker")
2428 {
2429 ready = true;
2430 break;
2431 }
2432 tokio::time::sleep(Duration::from_millis(30)).await;
2433 }
2434 assert!(ready, "worker never joined the pool");
2435
2436 let mut sched = HashMap::new();
2439 sched.insert(
2440 "mac-mini-monitor".to_string(),
2441 ResolvedSchedulablePlacement {
2442 pool: "mac-mini-monitor".into(),
2443 host_label: Some("mini".into()),
2444 },
2445 );
2446 let runner = bogus_sched_runner(HashMap::new(), sched).with_bus(Some(
2447 bamboo_subagent::BusEndpoint {
2448 endpoint: endpoint.clone(),
2449 token: "t".into(),
2450 },
2451 ));
2452
2453 let mut session = session_of_role("mac-mini-monitor", "hello scheduled");
2454 let job = job_for("child-1");
2455 let (event_tx, _rx) = mpsc::channel::<AgentEvent>(64);
2456 let cancel = CancellationToken::new();
2457
2458 tokio::time::timeout(
2459 Duration::from_secs(10),
2460 runner.execute_external_child(&mut session, &job, event_tx, cancel),
2461 )
2462 .await
2463 .expect("run did not hang")
2464 .expect("schedulable run succeeded over the bus (no local spawn)");
2465
2466 let last = session
2468 .messages
2469 .iter()
2470 .rev()
2471 .find(|m| matches!(m.role, Role::Assistant))
2472 .expect("an assistant reply was written back");
2473 assert!(last.content.contains("echo:"), "got {:?}", last.content);
2474
2475 let placement = session
2477 .metadata
2478 .get("placement")
2479 .expect("scheduled child session stamped with a placement");
2480 assert!(placement.contains(r#""kind":"remote""#), "{placement}");
2481 assert!(placement.contains(r#""host":"mini""#), "{placement}");
2482
2483 worker.abort();
2484 }
2485}