1use super::init::{
2 build_connect_manager, build_provider_handles, build_schedule_manager, build_spawn_scheduler,
3 init_mcp_manager, init_metrics_service, init_schedule_store, init_skill_manager, init_storage,
4 load_permission_checker, spawn_session_map_cleanup_task,
5};
6use super::tools::{build_base_tools, build_root_tools};
7use super::*;
8use crate::tools::OptionalSubagentModelResolver;
9use bamboo_agent_core::storage::Storage;
10
11impl AppState {
12 pub async fn new(bamboo_home_dir: PathBuf) -> Result<Self, AppError> {
54 bamboo_config::paths::init_bamboo_dir(bamboo_home_dir.clone());
57
58 let (config, config_facade) = match bamboo_config::ConfigFacade::open_or_migrate(
65 &bamboo_home_dir,
66 ) {
67 Ok(facade) => {
68 let facade = Arc::new(facade);
69 let config =
70 super::config_runtime::load_facade_effective_config(&facade, &bamboo_home_dir);
71 (config, Some(facade))
72 }
73 Err(error) => {
74 if bamboo_config::section_layout_is_active(&bamboo_home_dir).unwrap_or(false) {
75 return Err(AppError::InternalError(anyhow::anyhow!(
76 "modular configuration authority is unavailable: {error}"
77 )));
78 }
79 tracing::warn!(
80 error = %error,
81 "modular configuration facade is unavailable; retaining recovered legacy authority"
82 );
83 (
84 Config::from_data_dir_without_publish(Some(bamboo_home_dir.clone())),
85 None,
86 )
87 }
88 };
89 config.publish_env_vars();
90
91 if config.plugin_trust.enforcement_is_off() {
103 super::config_runtime::warn_plugin_trust_enforcement_off();
104 }
105
106 let provider_registry =
107 match bamboo_llm::ProviderRegistry::from_config(&config, bamboo_home_dir.clone()).await
108 {
109 Ok(registry) => Arc::new(registry),
110 Err(e) => {
111 tracing::error!("Failed to create provider registry: {}", e);
112 Arc::new(
113 bamboo_llm::ProviderRegistry::from_config(
114 &Config::default(),
115 bamboo_home_dir.clone(),
116 )
117 .await
118 .expect("Cannot create even an empty provider registry"),
119 )
120 }
121 };
122
123 let provider = provider_registry.get_default().unwrap_or_else(|| {
124 let default_provider_name = provider_registry.default_provider_name();
125 let message = if config.has_provider_instances() {
126 format!(
127 "Default provider instance '{}' is not available or failed to initialize",
128 default_provider_name
129 )
130 } else {
131 format!(
132 "Provider '{}' is not available or failed to initialize",
133 config.provider
134 )
135 };
136 Arc::new(UnconfiguredProvider { message }) as Arc<dyn LLMProvider>
137 });
138
139 Self::new_with_provider_and_facade(bamboo_home_dir, config, provider, config_facade).await
140 }
141
142 pub async fn new_with_provider(
157 bamboo_home_dir: PathBuf,
158 config: Config,
159 provider: Arc<dyn LLMProvider>,
160 ) -> Result<Self, AppError> {
161 Self::new_with_provider_and_facade(bamboo_home_dir, config, provider, None).await
162 }
163
164 async fn new_with_provider_and_facade(
165 bamboo_home_dir: PathBuf,
166 config: Config,
167 provider: Arc<dyn LLMProvider>,
168 config_facade: Option<Arc<bamboo_config::ConfigFacade>>,
169 ) -> Result<Self, AppError> {
170 let data_dir = bamboo_home_dir.clone();
172 let (session_store, storage) = init_storage(&data_dir).await?;
173 let project_store = Arc::new(bamboo_projects::ProjectStore::open(&data_dir).map_err(
174 |error| {
175 AppError::InternalError(anyhow::anyhow!(
176 "failed to initialize Project registry: {error}"
177 ))
178 },
179 )?);
180 let persistence = Arc::new(LockedSessionStore::new(storage.clone()));
181 let session_inbox: Arc<dyn bamboo_domain::SessionInboxPort> =
182 Arc::new(bamboo_storage::FileSessionInbox::new(
183 session_store.clone(),
184 bamboo_domain::SessionInboxLimits::default(),
185 ));
186 let session_activation_router = bamboo_engine::SessionActivationRouter::new();
187 let session_messenger = Arc::new(bamboo_engine::SessionMessenger::new(
188 storage.clone(),
189 session_inbox.clone(),
190 session_activation_router.clone(),
191 ));
192
193 let sessions: bamboo_engine::SessionCache = Arc::new(dashmap::DashMap::new());
195
196 let mut config = config;
203 let embedded_broker = maybe_embed_broker(&mut config, &data_dir).await;
204
205 let config = Arc::new(RwLock::new(config));
206
207 let live_default_workspace: Arc<dyn Fn() -> Option<PathBuf> + Send + Sync> = {
218 let config_for_workspace = config.clone();
219 let last_known: Arc<std::sync::Mutex<Option<PathBuf>>> =
220 Arc::new(std::sync::Mutex::new(None));
221 Arc::new(move || match config_for_workspace.try_read() {
222 Ok(cfg) => {
223 let path = cfg.get_default_work_area_path();
224 if let Ok(mut cache) = last_known.lock() {
225 *cache = path.clone();
226 }
227 path
228 }
229 Err(_) => last_known.lock().ok().and_then(|cache| cache.clone()),
230 })
231 };
232
233 let live_workspace_root: Arc<
238 dyn Fn() -> bamboo_agent_core::workspace_state::WorkspaceRootConfig + Send + Sync,
239 > = {
240 let app_data_dir = data_dir.clone();
241 Arc::new(
242 move || bamboo_agent_core::workspace_state::WorkspaceRootConfig {
243 root: bamboo_config::paths::resolve_workspace_root_in(&app_data_dir),
244 confine: bamboo_config::paths::workspace_confinement_enforced(),
245 },
246 )
247 };
248
249 let workspace_resolver = bamboo_agent_core::workspace_state::WorkspaceResolver::new(
250 {
251 let provider = live_default_workspace.clone();
252 move || provider()
253 },
254 {
255 let provider = live_workspace_root.clone();
256 move || provider()
257 },
258 );
259
260 bamboo_agent_core::workspace_state::set_default_workspace_provider(Box::new({
261 let provider = live_default_workspace;
262 move || provider()
263 }));
264 bamboo_agent_core::workspace_state::set_workspace_root_provider(Box::new({
265 let provider = live_workspace_root;
266 move || provider()
267 }));
268
269 let (permission_checker, permission_section) =
270 load_permission_checker(&bamboo_home_dir).await?;
271 let permission_io_lock = Arc::new(tokio::sync::Mutex::new(()));
272 let notification_service = Arc::new(bamboo_notification::NotificationService::new(
273 bamboo_home_dir.join("notification_preferences.json"),
274 ));
275 let session_watchers = super::watchers::SessionWatchers::new();
276 let (mcp_manager, _legacy_mcp_bootstrap) =
277 init_mcp_manager(config.clone(), &bamboo_home_dir);
278 let skill_manager = init_skill_manager(&data_dir).await;
279 let metrics_service = init_metrics_service(&data_dir).await?;
280
281 let startup_sessions = {
282 let entries = session_store.list_index_entries().await;
283 let mut sessions = Vec::new();
284 for entry in entries {
285 if let Some(session) = session_store
286 .load_session(&entry.id)
287 .await
288 .map_err(AppError::StorageError)?
289 {
290 sessions.push(session);
291 }
292 }
293 sessions
294 };
295 metrics_service
296 .reconcile_startup_sessions(startup_sessions, &[])
297 .await
298 .map_err(|error| {
299 AppError::InternalError(anyhow::anyhow!(
300 "Failed to reconcile stale metrics state on startup: {error}"
301 ))
302 })?;
303
304 let agent_runners: Arc<RwLock<HashMap<String, AgentRunner>>> =
305 Arc::new(RwLock::new(HashMap::new()));
306 let process_registry = Arc::new(ProcessRegistry::new());
311 let (provider_lock, provider_handle) = build_provider_handles(provider);
312
313 let config_snapshot = config.read().await;
315 let provider_registry = match bamboo_llm::ProviderRegistry::from_config(
316 &config_snapshot,
317 bamboo_home_dir.clone(),
318 )
319 .await
320 {
321 Ok(registry) => Arc::new(registry),
322 Err(e) => {
323 tracing::error!("Failed to create provider registry: {}", e);
324 Arc::new(
325 bamboo_llm::ProviderRegistry::from_config(
326 &Config::default(),
327 bamboo_home_dir.clone(),
328 )
329 .await
330 .expect("Cannot create even an empty provider registry"),
331 )
332 }
333 };
334 drop(config_snapshot);
335
336 let provider_router = Arc::new(bamboo_llm::ProviderModelRouter::new(
337 provider_registry.clone(),
338 ));
339 let model_catalog = Arc::new(bamboo_llm::ModelCatalogService::new(
340 provider_registry.clone(),
341 ));
342
343 let session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>> =
348 Arc::new(RwLock::new(HashMap::new()));
349
350 let notification_relay_deps = crate::app_state::session_events::NotificationRelayDeps {
356 notification_service: notification_service.clone(),
357 session_event_senders: session_event_senders.clone(),
358 session_watchers: session_watchers.clone(),
359 config: config.clone(),
360 };
361
362 let ledger_schedule_bridge =
365 Arc::new(crate::schedule_app::LateBoundLedgerBridge::default());
366
367 let session_repo = bamboo_engine::SessionRepository::new(
371 sessions.clone(),
372 storage.clone(),
373 persistence.clone(),
374 );
375
376 let account_sink = bamboo_engine::events::AccountEventSink::new(data_dir.join("events"))
380 .map_err(|e| {
381 AppError::InternalError(anyhow::anyhow!(
382 "failed to initialize account change-feed journal: {e}"
383 ))
384 })?;
385 let project_resource_watcher = super::project_watcher::ProjectResourceWatcher::start(
386 project_store.clone(),
387 account_sink.clone(),
388 std::time::Duration::from_millis(120),
389 )
390 .map_err(|error| {
391 AppError::InternalError(anyhow::anyhow!(
392 "failed to start Project resource watcher: {error}"
393 ))
394 })?;
395
396 let base_tools = build_base_tools(
397 config.clone(),
398 permission_checker.clone(),
399 mcp_manager.clone(),
400 skill_manager.clone(),
401 session_repo.clone(),
402 bamboo_home_dir.clone(),
403 notification_service.clone(),
404 session_event_senders.clone(),
405 session_watchers.clone(),
406 ledger_schedule_bridge.clone(),
407 project_store.clone(),
408 account_sink.clone(),
409 workspace_resolver.clone(),
410 );
411
412 let workflow_runs = crate::workflow::WorkflowRunAccess::new_with_permission_config(
416 &data_dir,
417 base_tools.clone(),
418 skill_manager.clone(),
419 session_repo.clone(),
420 permission_checker.permission_config(),
421 )
422 .await
423 .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?;
424
425 spawn_session_map_cleanup_task(agent_runners.clone(), session_event_senders.clone(), None);
429
430 {
433 let mut workflow_events = skill_manager.store().subscribe_workflow_catalog();
434 let account_sink = account_sink.clone();
435 tokio::spawn(async move {
436 loop {
437 let event = match workflow_events.recv().await {
438 Ok(event) => event,
439 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
440 tracing::warn!("Workflow catalog event bridge lagged by {skipped}");
441 continue;
442 }
443 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
444 };
445 if !event.public_workflow {
446 continue;
447 }
448 let event = match event.kind {
449 bamboo_skills::WorkflowCatalogEventKind::Changed => {
450 AgentEvent::WorkflowChanged {
451 workflow_id: event.workflow_id,
452 revision: event.revision,
453 scope: event.scope,
454 }
455 }
456 bamboo_skills::WorkflowCatalogEventKind::Invalid => {
457 AgentEvent::WorkflowInvalid {
458 workflow_id: event.workflow_id,
459 revision: event.revision,
460 scope: event.scope,
461 }
462 }
463 bamboo_skills::WorkflowCatalogEventKind::Recovered => {
464 AgentEvent::WorkflowRecovered {
465 workflow_id: event.workflow_id,
466 revision: event.revision,
467 scope: event.scope,
468 }
469 }
470 };
471 account_sink.record(None, &event);
472 }
473 });
474 }
475 let (approval_registry, restart_approval_events) =
476 bamboo_engine::external_agents::live::initialize_durable_approvals(
477 data_dir.join("approvals/child-approvals-v1.json"),
478 )
479 .map_err(|error| {
480 AppError::InternalError(anyhow::anyhow!(
481 "failed to initialize durable child approvals: {error}"
482 ))
483 })?;
484 for event in restart_approval_events {
485 account_sink.record(event.session_id(), &event);
486 }
487
488 let child_tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = base_tools.clone();
491
492 let project_context_resolver = Arc::new(
497 bamboo_engine::project_context::ProjectContextResolver::new_with_workspace_resolver(
498 Arc::new(crate::project_context::ProjectStoreContextSource::new(
499 project_store.clone(),
500 )),
501 workspace_resolver.clone(),
502 ),
503 );
504 let mut agent_builder = bamboo_engine::Agent::builder()
505 .storage(storage.clone())
506 .persistence(Arc::new(session_repo.clone()))
507 .session_inbox(session_inbox.clone())
508 .activation_router(session_activation_router.clone())
509 .session_messenger(session_messenger.clone())
510 .attachment_reader(session_store.clone())
511 .skill_manager(skill_manager.clone())
512 .metrics_collector(metrics_service.collector())
513 .config(config.clone())
514 .provider(provider_handle.clone())
515 .default_tools(base_tools.clone())
516 .project_context_resolver(project_context_resolver.clone());
517 if let Some(permission_config) = permission_checker.permission_config() {
518 agent_builder = agent_builder.permission_config(permission_config);
519 }
520 let agent = Arc::new(
521 agent_builder
522 .build()
523 .expect("agent runtime should be fully configured"),
524 );
525
526 let child_completion_coordinator =
527 Arc::new(bamboo_engine::ChildCompletionCoordinator::new(
528 storage.clone(),
529 persistence.clone(),
530 sessions.clone(),
531 agent_runners.clone(),
532 session_event_senders.clone(),
533 agent.clone(),
534 config.clone(),
535 provider_registry.clone(),
536 provider_router.clone(),
537 data_dir.clone(),
538 Some(account_sink.inbox()),
539 ));
540 session_activation_router
541 .set_spawner(child_completion_coordinator.clone())
542 .await;
543
544 let config_snapshot = config.read().await.clone();
546
547 let mcp_proxy_shutdown = tokio_util::sync::CancellationToken::new();
554 if let Some(broker) = config_snapshot.subagents().broker.clone() {
555 if !broker.endpoint.trim().is_empty() {
556 let backend: std::sync::Arc<dyn bamboo_agent_core::tools::ToolExecutor> =
557 std::sync::Arc::new(bamboo_mcp::executor::McpToolExecutor::new(
558 mcp_manager.clone(),
559 mcp_manager.tool_index(),
560 ));
561 let shutdown = mcp_proxy_shutdown.clone();
562
563 let role_entries: Vec<(String, Vec<String>)> = config_snapshot
572 .subagents()
573 .mcp_role_allowlist
574 .iter()
575 .map(|e| (e.role.clone(), e.tools.clone()))
576 .collect();
577 if role_entries.is_empty() {
578 tracing::info!(
584 "mcp proxy: no subagents.mcp_role_allowlist configured — every worker \
585 role sees/can call the full host-bound MCP tool set (opt in a role \
586 policy in config.json to scope tools per role; see issue #54)"
587 );
588 }
589 let known_tools: std::collections::HashSet<String> = backend
597 .list_tools()
598 .into_iter()
599 .map(|t| t.function.name)
600 .collect();
601 if !role_entries.is_empty() && known_tools.is_empty() {
602 tracing::warn!(
603 "mcp role allowlist: the orchestrator's MCP tool set was empty at \
604 policy-load time (servers may still be connecting in the background) — \
605 skipped tool-name typo validation for subagents.mcp_role_allowlist"
606 );
607 }
608 let allowlist = std::sync::Arc::new(bamboo_broker::RoleToolAllowlist::from_config(
609 role_entries,
610 &known_tools,
611 ));
612 tokio::spawn(async move {
613 let me = bamboo_broker::AgentRef {
614 session_id: bamboo_broker::ORCHESTRATOR_ID.to_string(),
615 role: Some("orchestrator".into()),
616 };
617 bamboo_broker::serve_mcp_proxy_supervised(
618 &broker.endpoint,
619 me,
620 &broker.token,
621 backend,
622 allowlist,
623 shutdown,
624 )
625 .await;
626 });
627 }
628 }
629 let parent_approval_reviewer = Arc::new(
630 crate::app_state::parent_approval_reviewer::ParentAgentApprovalReviewer::new(
631 session_repo.clone(),
632 provider_router.clone(),
633 ),
634 );
635 let codex_run_tokens = Arc::new(crate::codex_run_tokens::CodexRunTokenRegistry::default());
636 let external_runner =
637 bamboo_engine::external_agents::runtime::build_external_child_runner_with_codex_tokens(
638 &config_snapshot,
639 Some(approval_registry.clone()),
640 Some(parent_approval_reviewer),
641 permission_checker.permission_config(),
642 Some(codex_run_tokens.clone()),
643 );
644 external_runner.set_session_inbox_runtime(Some(
645 bamboo_engine::execution::spawn::SessionInboxRuntimeBinding {
646 router: session_activation_router.clone(),
647 inbox: session_inbox.clone(),
648 storage: storage.clone(),
649 persistence: persistence.clone(),
650 },
651 ));
652 let spawn_scheduler = build_spawn_scheduler(
653 agent.clone(),
654 child_tools,
655 sessions.clone(),
656 agent_runners.clone(),
657 session_event_senders.clone(),
658 external_runner,
659 Some(provider_router.clone()),
660 Some(child_completion_coordinator.clone()),
661 Some(data_dir.clone()),
662 Some(account_sink.inbox()),
663 Some(Arc::new(
664 crate::app_state::session_events::NotificationRelayLaunchHook::new(
665 notification_relay_deps.clone(),
666 ),
667 )),
668 );
669 child_completion_coordinator
670 .set_spawn_scheduler(&spawn_scheduler)
671 .await;
672
673 let tools_with_task = base_tools.clone();
674
675 let schedule_store = init_schedule_store(&data_dir).await?;
676
677 ledger_schedule_bridge
679 .bind(Arc::new(crate::schedule_app::ScheduleLedgerBridge::new(
680 schedule_store.clone(),
681 )))
682 .await;
683
684 let schedule_manager = build_schedule_manager(
685 schedule_store.clone(),
686 agent.clone(),
687 tools_with_task.clone(),
688 permission_checker.permission_config(),
689 sessions.clone(),
690 agent_runners.clone(),
691 session_event_senders.clone(),
692 persistence.clone(),
693 config.clone(),
694 provider_registry.clone(),
695 Some(data_dir.clone()),
696 Some(account_sink.inbox()),
697 notification_relay_deps.clone(),
698 project_store.clone(),
699 workspace_resolver.clone(),
700 );
701
702 bamboo_engine::auto_dream::spawn_auto_dream_task_with_project_resolver(
703 bamboo_engine::auto_dream::AutoDreamContext {
704 session_store: session_store.clone(),
705 storage: storage.clone(),
706 provider: provider_handle.clone(),
707 config: config.clone(),
708 provider_registry: provider_registry.clone(),
709 },
710 project_context_resolver.as_ref().clone(),
711 );
712
713 bamboo_engine::gardener::spawn_gardener_task_with_project_resolver(
717 bamboo_engine::auto_dream::AutoDreamContext {
718 session_store: session_store.clone(),
719 storage: storage.clone(),
720 provider: provider_handle.clone(),
721 config: config.clone(),
722 provider_registry: provider_registry.clone(),
723 },
724 project_context_resolver.clone(),
725 );
726
727 bamboo_engine::ledger_gardener::spawn_ledger_gardener_task(
731 bamboo_engine::ledger_gardener::LedgerGardenerContext {
732 dream: bamboo_engine::auto_dream::AutoDreamContext {
733 session_store: session_store.clone(),
734 storage: storage.clone(),
735 provider: provider_handle.clone(),
736 config: config.clone(),
737 provider_registry: provider_registry.clone(),
738 },
739 schedule_bridge: Some(ledger_schedule_bridge.clone()),
740 },
741 );
742
743 let config_for_resolver = config.clone();
744 let subagent_model_resolver: OptionalSubagentModelResolver = {
745 let registry = provider_registry.clone();
746 Some(Arc::new(
747 move |subagent_type: String| -> futures::future::BoxFuture<
748 'static,
749 Option<bamboo_domain::ProviderModelRef>,
750 > {
751 let config_for_resolver = config_for_resolver.clone();
752 let registry = registry.clone();
753 Box::pin(async move {
754 let config_snap = config_for_resolver.read().await.clone();
755 bamboo_engine::model_config_helper::resolve_subagent_model_ref(
756 &config_snap,
757 &config_snap.provider,
758 ®istry,
759 &subagent_type,
760 )
761 })
762 },
763 ))
764 };
765
766 let config_io_lock = Arc::new(tokio::sync::Mutex::new(()));
770 let fabric_registry: crate::tools::DeployedRegistry =
771 Arc::new(tokio::sync::Mutex::new(HashMap::new()));
772 let fabric_bamboo_bin =
773 std::env::current_exe().unwrap_or_else(|_| std::path::PathBuf::from("bamboo"));
774 let credential_store = Arc::new(bamboo_config::CredentialStore::open(&bamboo_home_dir));
775 let mut fabric_deployer = bamboo_server_tools::FabricDeployer::new(
776 config.clone(),
777 config_io_lock.clone(),
778 bamboo_home_dir.clone(),
779 fabric_registry,
780 fabric_bamboo_bin,
781 );
782 if let Some(facade) = config_facade.clone() {
783 let event_sink = account_sink.clone();
784 fabric_deployer = fabric_deployer.with_modular_persistence(
785 facade,
786 credential_store.clone(),
787 Arc::new(move |event| {
788 super::config_runtime::publish_registry_event(&event_sink, event);
789 }),
790 );
791 }
792 if let Err(error) = fabric_deployer.reconcile_stale_nodes_on_boot().await {
793 tracing::warn!(
794 error = %error,
795 "cluster-fabric boot reconcile failed; retaining the last adopted state"
796 );
797 }
798 let fabric_deployer = Arc::new(fabric_deployer);
799 let health_monitor = fabric_deployer
804 .clone()
805 .spawn_health_monitor()
806 .await
807 .map(HealthMonitor);
808
809 let tools = build_root_tools(
810 tools_with_task.clone(),
811 schedule_store.clone(),
812 schedule_manager.clone(),
813 session_store.clone(),
814 storage.clone(),
815 persistence.clone(),
816 session_messenger.clone(),
817 spawn_scheduler.clone(),
818 sessions.clone(),
819 agent_runners.clone(),
820 session_event_senders.clone(),
821 subagent_model_resolver,
822 config.clone(),
823 provider_registry.clone(),
824 config_snapshot.subagents().broker.clone(),
825 fabric_deployer.clone(),
826 project_store.clone(),
827 workspace_resolver.clone(),
828 );
829 let workflow_run_tool =
830 Arc::new(crate::workflow::WorkflowRunTool::new(workflow_runs.clone()));
831 let tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = Arc::new(
832 crate::tools::OverlayToolExecutor::new(tools, workflow_run_tool),
833 );
834
835 child_completion_coordinator
836 .set_root_tools(tools.clone())
837 .await;
838
839 for entry in session_store.list_index_entries().await {
844 match session_inbox.inspect(&entry.id).await {
845 Ok(backlog) if backlog.activation_pending() => {
846 if let Err(error) = bamboo_domain::SessionActivationPort::request_activation(
847 session_activation_router.as_ref(),
848 &entry.id,
849 backlog.activation_generation,
850 )
851 .await
852 {
853 tracing::error!(
854 session_id = %entry.id,
855 %error,
856 "failed to reactivate durable SessionInbox backlog during startup"
857 );
858 }
859 }
860 Ok(_) => {}
861 Err(error) => tracing::warn!(
862 session_id = %entry.id,
863 %error,
864 "failed to inspect SessionInbox during startup recovery"
865 ),
866 }
867 }
868
869 let tool_factory =
870 crate::tools::ToolSurfaceFactory::new(base_tools, tools_with_task, tools);
871
872 let session_repo = bamboo_engine::SessionRepository::new(
873 sessions.clone(),
874 storage.clone(),
875 persistence.clone(),
876 );
877
878 let connect_manager = Arc::new(
883 build_connect_manager(
884 agent.clone(),
885 tool_factory.get(crate::tools::ToolSurface::Root),
886 session_repo.clone(),
887 agent_runners.clone(),
888 session_event_senders.clone(),
889 Some(account_sink.inbox()),
890 Some(data_dir.clone()),
891 config.clone(),
892 provider_registry.clone(),
893 permission_checker.clone(),
894 project_store.clone(),
895 workspace_resolver.clone(),
896 )
897 .await
898 .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?,
899 );
900
901 let child_adapter = Arc::new(crate::tools::ChildSessionAdapter {
908 session_store: session_store.clone(),
909 storage: storage.clone(),
910 persistence: persistence.clone(),
911 session_messenger: Some(session_messenger.clone()),
912 scheduler: spawn_scheduler.clone(),
913 sessions_cache: sessions.clone(),
914 agent_runners: agent_runners.clone(),
915 session_event_senders: session_event_senders.clone(),
916 subagent_model_resolver: None,
917 config: config.clone(),
918 project_store: Some(project_store.clone()),
919 workspace_resolver: workspace_resolver.clone(),
920 parent_wait_slots: Arc::new(dashmap::DashMap::new()),
921 });
922 let guardian_spawner: Arc<dyn bamboo_engine::GuardianSpawner> = child_adapter.clone();
923 child_completion_coordinator
926 .set_guardian_spawner(guardian_spawner.clone())
927 .await;
928
929 let bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook> =
933 child_completion_coordinator.clone();
934
935 child_completion_coordinator.spawn_child_wait_watchdog();
941
942 let service_manager = Arc::new(crate::service_manager::ServiceManager::new());
948 let boot_reconcile_services_handle = {
960 let service_manager = service_manager.clone();
961 let app_data_dir = bamboo_home_dir.clone();
962 tokio::spawn(async move {
963 crate::plugin_installer::boot_reconcile_services(&app_data_dir, &service_manager)
964 .await;
965 })
966 };
967
968 let (config_watcher, config_live_health, mcp_config_live_health) =
969 super::config_runtime::ConfigWatcherRuntime::start(
970 bamboo_home_dir.clone(),
971 config.clone(),
972 config_facade.clone(),
973 config_io_lock.clone(),
974 provider_registry.clone(),
975 provider_lock.clone(),
976 mcp_manager.clone(),
977 account_sink.clone(),
978 );
979 Ok(Self {
980 app_data_dir: bamboo_home_dir,
981 config,
982 config_facade,
983 config_io_lock,
984 config_live_health,
985 mcp_config_live_health,
986 config_watcher,
987 project_resource_watcher,
988 credential_store,
989 fabric_deployer,
990 embedded_broker,
991 health_monitor,
992 provider: provider_lock,
993 provider_handle,
994 sessions,
995 storage,
996 session_store,
997 project_store,
998 project_context_resolver,
999 workspace_resolver,
1000 session_repo,
1001 persistence,
1002 session_inbox,
1003 session_activation_router,
1004 session_messenger,
1005 spawn_scheduler,
1006 child_completion_coordinator,
1007 guardian_spawner,
1008 bash_resume_hook,
1009 schedule_store,
1010 schedule_manager,
1011 connect_manager,
1012 tool_factory,
1013 permission_checker,
1014 permission_section,
1015 permission_io_lock,
1016 approval_registry,
1017 notification_service,
1018 session_watchers,
1019 cancel_tokens: Arc::new(RwLock::new(HashMap::new())),
1020 mcp_proxy_shutdown,
1021 skill_manager,
1022 workflow_runs,
1023 mcp_manager,
1024 service_manager,
1025 boot_reconcile_services_handle: tokio::sync::Mutex::new(Some(
1026 boot_reconcile_services_handle,
1027 )),
1028 metrics_service,
1029 agent_runners,
1030 execute_startups: Arc::new(std::sync::Mutex::new(HashMap::new())),
1031 session_event_senders,
1032 account_sink,
1033 process_registry,
1034 metrics_bus: None, agent,
1036 provider_registry,
1037 provider_router,
1038 model_catalog,
1039 title_gen_in_flight: Arc::new(dashmap::DashSet::new()),
1040 pairing_codes: Arc::new(dashmap::DashMap::new()),
1041 pairing_code_guard: Arc::new(crate::handlers::settings::PairingCodeGuard::default()),
1042 root_password_guard: Arc::new(crate::handlers::settings::RootPasswordGuard::default()),
1043 codex_run_tokens,
1044 })
1046 }
1047}
1048
1049pub struct EmbeddedBroker {
1052 task: tokio::task::JoinHandle<()>,
1053 gc_task: tokio::task::JoinHandle<()>,
1054}
1055
1056impl Drop for EmbeddedBroker {
1057 fn drop(&mut self) {
1058 self.task.abort();
1059 self.gc_task.abort();
1060 }
1061}
1062
1063pub struct HealthMonitor(tokio::task::JoinHandle<()>);
1067
1068impl Drop for HealthMonitor {
1069 fn drop(&mut self) {
1070 self.0.abort();
1071 }
1072}
1073
1074async fn maybe_embed_broker(
1086 config: &mut bamboo_llm::Config,
1087 data_dir: &std::path::Path,
1088) -> Option<EmbeddedBroker> {
1089 if let Some(external) = load_external_broker(data_dir) {
1094 let endpoint = external.endpoint.trim().to_string();
1095 if broker_endpoint_reachable(&endpoint).await {
1096 tracing::info!(%endpoint, "using external broker from broker.json");
1097 config.subagents_mut().broker = Some(external);
1098 return None;
1099 }
1100 tracing::warn!(
1101 %endpoint,
1102 "broker.json endpoint is unreachable — embedding a fresh in-process broker instead"
1103 );
1104 }
1105
1106 let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
1107 Ok(l) => l,
1108 Err(e) => {
1109 tracing::warn!("embedded broker: bind failed, sub-agent dispatch disabled: {e}");
1110 return None;
1111 }
1112 };
1113 let port = match listener.local_addr() {
1114 Ok(a) => a.port(),
1115 Err(e) => {
1116 tracing::warn!("embedded broker: local_addr failed: {e}");
1117 return None;
1118 }
1119 };
1120 let token = uuid::Uuid::new_v4().simple().to_string();
1121 let root = data_dir.join("broker");
1122 let core = Arc::new(bamboo_broker::BrokerCore::new(root));
1123 let gc_task = core
1126 .clone()
1127 .spawn_mailbox_gc(std::time::Duration::from_secs(300));
1128 let server = Arc::new(bamboo_broker::BrokerServer::new(core, token.clone()));
1129
1130 let task = tokio::spawn(async move {
1131 if let Err(e) = server.serve(listener).await {
1132 tracing::error!("embedded broker serve loop ended: {e}");
1133 }
1134 });
1135
1136 config.subagents_mut().broker = Some(bamboo_config::BrokerClientConfig {
1139 endpoint: format!("ws://127.0.0.1:{port}"),
1140 token,
1141 token_encrypted: None,
1142 credential_ref: None,
1143 configured: false,
1144 });
1145 tracing::info!(port, "embedded mailbox bus (broker) started in-process");
1146 Some(EmbeddedBroker { task, gc_task })
1147}
1148
1149fn load_external_broker(data_dir: &std::path::Path) -> Option<bamboo_config::BrokerClientConfig> {
1158 let path = data_dir.join("broker.json");
1159 if let Err(error) = bamboo_config::migrate_external_broker_credentials(data_dir)
1160 .and_then(|_| bamboo_config::ensure_provider_mcp_migration_ready(data_dir))
1161 {
1162 tracing::warn!(error = %error, "external broker credential migration unavailable");
1163 return None;
1164 }
1165 let bytes = std::fs::read(&path).ok()?;
1166 match parse_external_broker_snapshot(&bytes, data_dir) {
1167 Ok(cfg) => Some(cfg),
1168 Err(error) => {
1169 tracing::warn!(?path, %error, "broker.json is unavailable");
1170 None
1171 }
1172 }
1173}
1174
1175fn parse_external_broker_snapshot(
1176 bytes: &[u8],
1177 data_dir: &std::path::Path,
1178) -> bamboo_config::ConfigStoreResult<bamboo_config::BrokerClientConfig> {
1179 let mut cfg: bamboo_config::BrokerClientConfig = serde_json::from_slice(bytes)?;
1180 if !cfg.token.trim().is_empty()
1181 || cfg
1182 .token_encrypted
1183 .as_deref()
1184 .is_some_and(|value| !value.trim().is_empty())
1185 {
1186 return Err(bamboo_config::ConfigStoreError::Validation(
1187 "legacy broker credential appeared after migration".to_string(),
1188 ));
1189 }
1190 if cfg.endpoint.trim().is_empty() {
1191 return Err(bamboo_config::ConfigStoreError::Validation(
1192 "broker endpoint is empty".to_string(),
1193 ));
1194 }
1195 cfg.hydrate_credential_from_store(data_dir)?;
1196 Ok(cfg)
1197}
1198
1199async fn broker_endpoint_reachable(endpoint: &str) -> bool {
1204 let host_port = endpoint
1205 .trim()
1206 .trim_start_matches("wss://")
1207 .trim_start_matches("ws://")
1208 .split('/')
1209 .next()
1210 .unwrap_or("");
1211 if host_port.is_empty() {
1212 return false;
1213 }
1214 matches!(
1215 tokio::time::timeout(
1216 std::time::Duration::from_millis(500),
1217 tokio::net::TcpStream::connect(host_port),
1218 )
1219 .await,
1220 Ok(Ok(_))
1221 )
1222}
1223
1224#[cfg(test)]
1225mod fabric_boot_reconcile_tests {
1226 use super::*;
1227 use bamboo_config::cluster_fabric::{
1228 DeployProfile, Node, NodePlacement, NodeState, NodeStatus, TrustLevel,
1229 };
1230 use bamboo_config::{ClusterNodeCredentialIntents, ConfigFacade};
1231 use std::collections::BTreeMap;
1232
1233 #[tokio::test]
1234 async fn restart_reconcile_keeps_runtime_process_facade_and_disk_on_one_revision() {
1235 let _key = bamboo_config::encryption::set_test_encryption_key([0x7b; 32]);
1236 let dir = tempfile::tempdir().unwrap();
1237 let first = AppState::new(dir.path().to_path_buf()).await.unwrap();
1238 first
1239 .update_cluster_fabric_credentials(
1240 0,
1241 BTreeMap::from([(
1242 "boot-node".to_string(),
1243 ClusterNodeCredentialIntents::clear_all(),
1244 )]),
1245 |config| {
1246 config.cluster_fabric.nodes.push(Node {
1247 id: "boot-node".to_string(),
1248 label: "boot-node".to_string(),
1249 placement: NodePlacement::Local,
1250 trust_level: TrustLevel::Trusted,
1251 deploy: DeployProfile::default(),
1252 state: Some(NodeState {
1253 status: NodeStatus::Running,
1254 worker_id: Some("stale-worker".to_string()),
1255 ..Default::default()
1256 }),
1257 enabled: true,
1258 });
1259 Ok(())
1260 },
1261 )
1262 .await
1263 .unwrap();
1264 assert_eq!(
1265 first
1266 .config_facade
1267 .as_ref()
1268 .unwrap()
1269 .registry()
1270 .cluster_fabric
1271 .snapshot()
1272 .revision,
1273 1
1274 );
1275 drop(first);
1276
1277 let restarted = AppState::new(dir.path().to_path_buf()).await.unwrap();
1278 let process_snapshot = restarted
1279 .config_facade
1280 .as_ref()
1281 .unwrap()
1282 .registry()
1283 .cluster_fabric
1284 .snapshot();
1285 assert_eq!(process_snapshot.revision, 2);
1286 assert_eq!(
1287 process_snapshot
1288 .data
1289 .0
1290 .node("boot-node")
1291 .unwrap()
1292 .state
1293 .as_ref()
1294 .unwrap()
1295 .status,
1296 NodeStatus::Unreachable
1297 );
1298 assert_eq!(
1299 restarted
1300 .config
1301 .read()
1302 .await
1303 .cluster_fabric
1304 .node("boot-node")
1305 .unwrap()
1306 .state
1307 .as_ref()
1308 .unwrap()
1309 .status,
1310 NodeStatus::Unreachable
1311 );
1312 let reopened = ConfigFacade::open(dir.path()).unwrap();
1313 assert_eq!(reopened.registry().cluster_fabric.snapshot().revision, 2);
1314 assert_eq!(
1315 reopened
1316 .effective_config()
1317 .cluster_fabric
1318 .node("boot-node")
1319 .unwrap()
1320 .state
1321 .as_ref()
1322 .unwrap()
1323 .status,
1324 NodeStatus::Unreachable
1325 );
1326 }
1327}
1328
1329#[cfg(test)]
1330mod broker_embed_tests {
1331 use super::{broker_endpoint_reachable, load_external_broker, parse_external_broker_snapshot};
1332
1333 #[test]
1334 fn load_external_broker_reads_broker_json_not_config() {
1335 let _key = bamboo_config::encryption::set_test_encryption_key([0x57; 32]);
1336 let dir = tempfile::tempdir().unwrap();
1337 assert!(load_external_broker(dir.path()).is_none());
1339
1340 std::fs::write(
1342 dir.path().join("broker.json"),
1343 r#"{ "endpoint": "wss://broker.example:9600", "token": "t" }"#,
1344 )
1345 .unwrap();
1346 let got = load_external_broker(dir.path()).expect("parsed");
1347 assert_eq!(got.endpoint, "wss://broker.example:9600");
1348 assert_eq!(got.token, "t");
1349 assert_eq!(
1350 got.credential_ref.as_ref().unwrap().as_str(),
1351 "broker.external.bearer_token"
1352 );
1353 assert!(got.configured);
1354 let durable = std::fs::read_to_string(dir.path().join("broker.json")).unwrap();
1355 assert!(!durable.contains("\"token\""));
1356 assert!(!durable.contains("token_encrypted"));
1357 assert!(!durable.contains("\"t\""));
1358
1359 std::fs::write(dir.path().join("broker.json"), r#"{ "endpoint": " " }"#).unwrap();
1361 assert!(load_external_broker(dir.path()).is_none());
1362
1363 std::fs::write(dir.path().join("broker.json"), "not json").unwrap();
1365 assert!(load_external_broker(dir.path()).is_none());
1366 }
1367
1368 #[test]
1369 fn external_broker_reference_fails_closed_and_tracks_generic_cas_updates() {
1370 let _key = bamboo_config::encryption::set_test_encryption_key([0x58; 32]);
1371 let dir = tempfile::tempdir().unwrap();
1372 let reference =
1373 bamboo_config::credential_ref("broker", "external", "bearer_token").unwrap();
1374 std::fs::write(
1375 dir.path().join("broker.json"),
1376 serde_json::to_vec_pretty(&serde_json::json!({
1377 "endpoint": "wss://broker.example:9600",
1378 "credential_ref": reference,
1379 "configured": true,
1380 }))
1381 .unwrap(),
1382 )
1383 .unwrap();
1384
1385 assert!(load_external_broker(dir.path()).is_none());
1386 std::fs::write(
1387 dir.path().join("broker.json"),
1388 serde_json::to_vec_pretty(&serde_json::json!({
1389 "endpoint": "wss://broker.example:9600",
1390 "credential_ref": reference,
1391 "configured": false,
1392 }))
1393 .unwrap(),
1394 )
1395 .unwrap();
1396 assert!(load_external_broker(dir.path()).is_none());
1397 let store = bamboo_config::CredentialStore::open(dir.path());
1398 store
1399 .replace(
1400 reference.clone(),
1401 "replacement-token",
1402 bamboo_config::CredentialSource::User,
1403 0,
1404 )
1405 .unwrap();
1406 let hydrated = load_external_broker(dir.path()).unwrap();
1407 assert_eq!(hydrated.token, "replacement-token");
1408 store.clear(&reference, 1).unwrap();
1409 assert!(load_external_broker(dir.path()).is_none());
1410 }
1411
1412 #[test]
1413 fn broker_snapshot_with_late_legacy_secret_fails_closed() {
1414 let _key = bamboo_config::encryption::set_test_encryption_key([0x59; 32]);
1415 let dir = tempfile::tempdir().unwrap();
1416 let error = parse_external_broker_snapshot(
1417 br#"{
1418 "endpoint": "wss://broker.example:9600",
1419 "token": "late-legacy-token"
1420 }"#,
1421 dir.path(),
1422 )
1423 .unwrap_err();
1424 assert!(error.to_string().contains("appeared after migration"));
1425 }
1426
1427 #[tokio::test]
1428 async fn reachability_probe_distinguishes_live_from_dead() {
1429 let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1432 let dead = l.local_addr().unwrap();
1433 drop(l);
1434 assert!(!broker_endpoint_reachable(&format!("ws://{dead}")).await);
1435
1436 assert!(!broker_endpoint_reachable("").await);
1438 assert!(!broker_endpoint_reachable("ws://").await);
1439
1440 let live = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1442 let addr = live.local_addr().unwrap();
1443 assert!(broker_endpoint_reachable(&format!("ws://{addr}/stream")).await);
1444 }
1445}