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 {
216 let config_for_workspace = config.clone();
217 let last_known: Arc<std::sync::Mutex<Option<PathBuf>>> =
218 Arc::new(std::sync::Mutex::new(None));
219 bamboo_agent_core::workspace_state::set_default_workspace_provider(Box::new(
220 move || match config_for_workspace.try_read() {
221 Ok(cfg) => {
222 let path = cfg.get_default_work_area_path();
223 if let Ok(mut cache) = last_known.lock() {
224 *cache = path.clone();
225 }
226 path
227 }
228 Err(_) => last_known.lock().ok().and_then(|c| c.clone()),
229 },
230 ));
231 }
232
233 bamboo_agent_core::workspace_state::set_workspace_root_provider(Box::new(|| {
244 bamboo_agent_core::workspace_state::WorkspaceRootConfig {
245 root: bamboo_config::paths::resolve_workspace_root(),
246 confine: bamboo_config::paths::workspace_confinement_enforced(),
247 }
248 }));
249
250 let (permission_checker, permission_section) =
251 load_permission_checker(&bamboo_home_dir).await?;
252 let permission_io_lock = Arc::new(tokio::sync::Mutex::new(()));
253 let notification_service = Arc::new(bamboo_notification::NotificationService::new(
254 bamboo_home_dir.join("notification_preferences.json"),
255 ));
256 let session_watchers = super::watchers::SessionWatchers::new();
257 let (mcp_manager, _legacy_mcp_bootstrap) =
258 init_mcp_manager(config.clone(), &bamboo_home_dir);
259 let skill_manager = init_skill_manager(&data_dir).await;
260 let metrics_service = init_metrics_service(&data_dir).await?;
261
262 let startup_sessions = {
263 let entries = session_store.list_index_entries().await;
264 let mut sessions = Vec::new();
265 for entry in entries {
266 if let Some(session) = session_store
267 .load_session(&entry.id)
268 .await
269 .map_err(AppError::StorageError)?
270 {
271 sessions.push(session);
272 }
273 }
274 sessions
275 };
276 metrics_service
277 .reconcile_startup_sessions(startup_sessions, &[])
278 .await
279 .map_err(|error| {
280 AppError::InternalError(anyhow::anyhow!(
281 "Failed to reconcile stale metrics state on startup: {error}"
282 ))
283 })?;
284
285 let agent_runners: Arc<RwLock<HashMap<String, AgentRunner>>> =
286 Arc::new(RwLock::new(HashMap::new()));
287 let process_registry = Arc::new(ProcessRegistry::new());
292 let (provider_lock, provider_handle) = build_provider_handles(provider);
293
294 let config_snapshot = config.read().await;
296 let provider_registry = match bamboo_llm::ProviderRegistry::from_config(
297 &config_snapshot,
298 bamboo_home_dir.clone(),
299 )
300 .await
301 {
302 Ok(registry) => Arc::new(registry),
303 Err(e) => {
304 tracing::error!("Failed to create provider registry: {}", e);
305 Arc::new(
306 bamboo_llm::ProviderRegistry::from_config(
307 &Config::default(),
308 bamboo_home_dir.clone(),
309 )
310 .await
311 .expect("Cannot create even an empty provider registry"),
312 )
313 }
314 };
315 drop(config_snapshot);
316
317 let provider_router = Arc::new(bamboo_llm::ProviderModelRouter::new(
318 provider_registry.clone(),
319 ));
320 let model_catalog = Arc::new(bamboo_llm::ModelCatalogService::new(
321 provider_registry.clone(),
322 ));
323
324 let session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>> =
329 Arc::new(RwLock::new(HashMap::new()));
330
331 let notification_relay_deps = crate::app_state::session_events::NotificationRelayDeps {
337 notification_service: notification_service.clone(),
338 session_event_senders: session_event_senders.clone(),
339 session_watchers: session_watchers.clone(),
340 config: config.clone(),
341 };
342
343 let ledger_schedule_bridge =
346 Arc::new(crate::schedule_app::LateBoundLedgerBridge::default());
347
348 let session_repo = bamboo_engine::SessionRepository::new(
352 sessions.clone(),
353 storage.clone(),
354 persistence.clone(),
355 );
356
357 let account_sink = bamboo_engine::events::AccountEventSink::new(data_dir.join("events"))
361 .map_err(|e| {
362 AppError::InternalError(anyhow::anyhow!(
363 "failed to initialize account change-feed journal: {e}"
364 ))
365 })?;
366 let project_resource_watcher = super::project_watcher::ProjectResourceWatcher::start(
367 project_store.clone(),
368 account_sink.clone(),
369 std::time::Duration::from_millis(120),
370 )
371 .map_err(|error| {
372 AppError::InternalError(anyhow::anyhow!(
373 "failed to start Project resource watcher: {error}"
374 ))
375 })?;
376
377 let base_tools = build_base_tools(
378 config.clone(),
379 permission_checker.clone(),
380 mcp_manager.clone(),
381 skill_manager.clone(),
382 session_repo.clone(),
383 bamboo_home_dir.clone(),
384 notification_service.clone(),
385 session_event_senders.clone(),
386 session_watchers.clone(),
387 ledger_schedule_bridge.clone(),
388 project_store.clone(),
389 account_sink.clone(),
390 );
391
392 let workflow_runs = crate::workflow::WorkflowRunAccess::new(
396 &data_dir,
397 base_tools.clone(),
398 skill_manager.clone(),
399 session_repo.clone(),
400 )
401 .await
402 .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?;
403
404 spawn_session_map_cleanup_task(agent_runners.clone(), session_event_senders.clone(), None);
408
409 {
412 let mut workflow_events = skill_manager.store().subscribe_workflow_catalog();
413 let account_sink = account_sink.clone();
414 tokio::spawn(async move {
415 loop {
416 let event = match workflow_events.recv().await {
417 Ok(event) => event,
418 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
419 tracing::warn!("Workflow catalog event bridge lagged by {skipped}");
420 continue;
421 }
422 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
423 };
424 let event = match event.kind {
425 bamboo_skills::WorkflowCatalogEventKind::Changed => {
426 AgentEvent::WorkflowChanged {
427 workflow_id: event.workflow_id,
428 revision: event.revision,
429 scope: event.scope,
430 }
431 }
432 bamboo_skills::WorkflowCatalogEventKind::Invalid => {
433 AgentEvent::WorkflowInvalid {
434 workflow_id: event.workflow_id,
435 revision: event.revision,
436 scope: event.scope,
437 }
438 }
439 bamboo_skills::WorkflowCatalogEventKind::Recovered => {
440 AgentEvent::WorkflowRecovered {
441 workflow_id: event.workflow_id,
442 revision: event.revision,
443 scope: event.scope,
444 }
445 }
446 };
447 account_sink.record(None, &event);
448 }
449 });
450 }
451 let (approval_registry, restart_approval_events) =
452 bamboo_engine::external_agents::live::initialize_durable_approvals(
453 data_dir.join("approvals/child-approvals-v1.json"),
454 )
455 .map_err(|error| {
456 AppError::InternalError(anyhow::anyhow!(
457 "failed to initialize durable child approvals: {error}"
458 ))
459 })?;
460 for event in restart_approval_events {
461 account_sink.record(event.session_id(), &event);
462 }
463
464 let child_tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = base_tools.clone();
467
468 let project_context_resolver = Arc::new(
473 bamboo_engine::project_context::ProjectContextResolver::new(Arc::new(
474 crate::project_context::ProjectStoreContextSource::new(project_store.clone()),
475 )),
476 );
477 let agent = Arc::new(
478 bamboo_engine::Agent::builder()
479 .storage(storage.clone())
480 .persistence(Arc::new(session_repo.clone()))
481 .session_inbox(session_inbox.clone())
482 .activation_router(session_activation_router.clone())
483 .session_messenger(session_messenger.clone())
484 .attachment_reader(session_store.clone())
485 .skill_manager(skill_manager.clone())
486 .metrics_collector(metrics_service.collector())
487 .config(config.clone())
488 .provider(provider_handle.clone())
489 .default_tools(base_tools.clone())
490 .project_context_resolver(project_context_resolver.clone())
491 .build()
492 .expect("agent runtime should be fully configured"),
493 );
494
495 let child_completion_coordinator =
496 Arc::new(bamboo_engine::ChildCompletionCoordinator::new(
497 storage.clone(),
498 persistence.clone(),
499 sessions.clone(),
500 agent_runners.clone(),
501 session_event_senders.clone(),
502 agent.clone(),
503 config.clone(),
504 provider_registry.clone(),
505 provider_router.clone(),
506 data_dir.clone(),
507 Some(account_sink.inbox()),
508 ));
509 session_activation_router
510 .set_spawner(child_completion_coordinator.clone())
511 .await;
512
513 let config_snapshot = config.read().await.clone();
515
516 let mcp_proxy_shutdown = tokio_util::sync::CancellationToken::new();
523 if let Some(broker) = config_snapshot.subagents().broker.clone() {
524 if !broker.endpoint.trim().is_empty() {
525 let backend: std::sync::Arc<dyn bamboo_agent_core::tools::ToolExecutor> =
526 std::sync::Arc::new(bamboo_mcp::executor::McpToolExecutor::new(
527 mcp_manager.clone(),
528 mcp_manager.tool_index(),
529 ));
530 let shutdown = mcp_proxy_shutdown.clone();
531
532 let role_entries: Vec<(String, Vec<String>)> = config_snapshot
541 .subagents()
542 .mcp_role_allowlist
543 .iter()
544 .map(|e| (e.role.clone(), e.tools.clone()))
545 .collect();
546 if role_entries.is_empty() {
547 tracing::info!(
553 "mcp proxy: no subagents.mcp_role_allowlist configured — every worker \
554 role sees/can call the full host-bound MCP tool set (opt in a role \
555 policy in config.json to scope tools per role; see issue #54)"
556 );
557 }
558 let known_tools: std::collections::HashSet<String> = backend
566 .list_tools()
567 .into_iter()
568 .map(|t| t.function.name)
569 .collect();
570 if !role_entries.is_empty() && known_tools.is_empty() {
571 tracing::warn!(
572 "mcp role allowlist: the orchestrator's MCP tool set was empty at \
573 policy-load time (servers may still be connecting in the background) — \
574 skipped tool-name typo validation for subagents.mcp_role_allowlist"
575 );
576 }
577 let allowlist = std::sync::Arc::new(bamboo_broker::RoleToolAllowlist::from_config(
578 role_entries,
579 &known_tools,
580 ));
581 tokio::spawn(async move {
582 let me = bamboo_broker::AgentRef {
583 session_id: bamboo_broker::ORCHESTRATOR_ID.to_string(),
584 role: Some("orchestrator".into()),
585 };
586 bamboo_broker::serve_mcp_proxy_supervised(
587 &broker.endpoint,
588 me,
589 &broker.token,
590 backend,
591 allowlist,
592 shutdown,
593 )
594 .await;
595 });
596 }
597 }
598 let parent_approval_reviewer = Arc::new(
599 crate::app_state::parent_approval_reviewer::ParentAgentApprovalReviewer::new(
600 session_repo.clone(),
601 provider_router.clone(),
602 ),
603 );
604 let codex_run_tokens = Arc::new(crate::codex_run_tokens::CodexRunTokenRegistry::default());
605 let external_runner =
606 bamboo_engine::external_agents::runtime::build_external_child_runner_with_codex_tokens(
607 &config_snapshot,
608 Some(approval_registry.clone()),
609 Some(parent_approval_reviewer),
610 permission_checker.permission_config(),
611 Some(codex_run_tokens.clone()),
612 );
613 external_runner.set_session_inbox_runtime(Some(
614 bamboo_engine::execution::spawn::SessionInboxRuntimeBinding {
615 router: session_activation_router.clone(),
616 inbox: session_inbox.clone(),
617 storage: storage.clone(),
618 persistence: persistence.clone(),
619 },
620 ));
621 let spawn_scheduler = build_spawn_scheduler(
622 agent.clone(),
623 child_tools,
624 sessions.clone(),
625 agent_runners.clone(),
626 session_event_senders.clone(),
627 external_runner,
628 Some(provider_router.clone()),
629 Some(child_completion_coordinator.clone()),
630 Some(data_dir.clone()),
631 Some(account_sink.inbox()),
632 Some(Arc::new(
633 crate::app_state::session_events::NotificationRelayLaunchHook::new(
634 notification_relay_deps.clone(),
635 ),
636 )),
637 );
638 child_completion_coordinator
639 .set_spawn_scheduler(&spawn_scheduler)
640 .await;
641
642 let tools_with_task = base_tools.clone();
643
644 let schedule_store = init_schedule_store(&data_dir).await?;
645
646 ledger_schedule_bridge
648 .bind(Arc::new(crate::schedule_app::ScheduleLedgerBridge::new(
649 schedule_store.clone(),
650 )))
651 .await;
652
653 let schedule_manager = build_schedule_manager(
654 schedule_store.clone(),
655 agent.clone(),
656 tools_with_task.clone(),
657 permission_checker.permission_config(),
658 sessions.clone(),
659 agent_runners.clone(),
660 session_event_senders.clone(),
661 persistence.clone(),
662 config.clone(),
663 provider_registry.clone(),
664 Some(data_dir.clone()),
665 Some(account_sink.inbox()),
666 notification_relay_deps.clone(),
667 project_store.clone(),
668 );
669
670 bamboo_engine::auto_dream::spawn_auto_dream_task_with_project_resolver(
671 bamboo_engine::auto_dream::AutoDreamContext {
672 session_store: session_store.clone(),
673 storage: storage.clone(),
674 provider: provider_handle.clone(),
675 config: config.clone(),
676 provider_registry: provider_registry.clone(),
677 },
678 project_context_resolver.as_ref().clone(),
679 );
680
681 bamboo_engine::gardener::spawn_gardener_task_with_project_resolver(
685 bamboo_engine::auto_dream::AutoDreamContext {
686 session_store: session_store.clone(),
687 storage: storage.clone(),
688 provider: provider_handle.clone(),
689 config: config.clone(),
690 provider_registry: provider_registry.clone(),
691 },
692 project_context_resolver.clone(),
693 );
694
695 bamboo_engine::ledger_gardener::spawn_ledger_gardener_task(
699 bamboo_engine::ledger_gardener::LedgerGardenerContext {
700 dream: bamboo_engine::auto_dream::AutoDreamContext {
701 session_store: session_store.clone(),
702 storage: storage.clone(),
703 provider: provider_handle.clone(),
704 config: config.clone(),
705 provider_registry: provider_registry.clone(),
706 },
707 schedule_bridge: Some(ledger_schedule_bridge.clone()),
708 },
709 );
710
711 let config_for_resolver = config.clone();
712 let subagent_model_resolver: OptionalSubagentModelResolver = {
713 let registry = provider_registry.clone();
714 Some(Arc::new(
715 move |subagent_type: String| -> futures::future::BoxFuture<
716 'static,
717 Option<bamboo_domain::ProviderModelRef>,
718 > {
719 let config_for_resolver = config_for_resolver.clone();
720 let registry = registry.clone();
721 Box::pin(async move {
722 let config_snap = config_for_resolver.read().await.clone();
723 bamboo_engine::model_config_helper::resolve_subagent_model_ref(
724 &config_snap,
725 &config_snap.provider,
726 ®istry,
727 &subagent_type,
728 )
729 })
730 },
731 ))
732 };
733
734 let config_io_lock = Arc::new(tokio::sync::Mutex::new(()));
738 let fabric_registry: crate::tools::DeployedRegistry =
739 Arc::new(tokio::sync::Mutex::new(HashMap::new()));
740 let fabric_bamboo_bin =
741 std::env::current_exe().unwrap_or_else(|_| std::path::PathBuf::from("bamboo"));
742 let fabric_deployer = Arc::new(bamboo_server_tools::FabricDeployer::new(
743 config.clone(),
744 config_io_lock.clone(),
745 bamboo_home_dir.clone(),
746 fabric_registry,
747 fabric_bamboo_bin,
748 ));
749 let health_monitor = fabric_deployer
754 .clone()
755 .spawn_health_monitor()
756 .await
757 .map(HealthMonitor);
758
759 let tools = build_root_tools(
760 tools_with_task.clone(),
761 schedule_store.clone(),
762 schedule_manager.clone(),
763 session_store.clone(),
764 storage.clone(),
765 persistence.clone(),
766 session_messenger.clone(),
767 spawn_scheduler.clone(),
768 sessions.clone(),
769 agent_runners.clone(),
770 session_event_senders.clone(),
771 subagent_model_resolver,
772 config.clone(),
773 provider_registry.clone(),
774 config_snapshot.subagents().broker.clone(),
775 fabric_deployer.clone(),
776 project_store.clone(),
777 );
778 let workflow_run_tool =
779 Arc::new(crate::workflow::WorkflowRunTool::new(workflow_runs.clone()));
780 let tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = Arc::new(
781 crate::tools::OverlayToolExecutor::new(tools, workflow_run_tool),
782 );
783
784 child_completion_coordinator
785 .set_root_tools(tools.clone())
786 .await;
787
788 for entry in session_store.list_index_entries().await {
793 match session_inbox.inspect(&entry.id).await {
794 Ok(backlog) if backlog.activation_pending() => {
795 if let Err(error) = bamboo_domain::SessionActivationPort::request_activation(
796 session_activation_router.as_ref(),
797 &entry.id,
798 backlog.activation_generation,
799 )
800 .await
801 {
802 tracing::error!(
803 session_id = %entry.id,
804 %error,
805 "failed to reactivate durable SessionInbox backlog during startup"
806 );
807 }
808 }
809 Ok(_) => {}
810 Err(error) => tracing::warn!(
811 session_id = %entry.id,
812 %error,
813 "failed to inspect SessionInbox during startup recovery"
814 ),
815 }
816 }
817
818 let tool_factory =
819 crate::tools::ToolSurfaceFactory::new(base_tools, tools_with_task, tools);
820
821 let session_repo = bamboo_engine::SessionRepository::new(
822 sessions.clone(),
823 storage.clone(),
824 persistence.clone(),
825 );
826
827 let connect_manager = Arc::new(
832 build_connect_manager(
833 agent.clone(),
834 tool_factory.get(crate::tools::ToolSurface::Root),
835 session_repo.clone(),
836 agent_runners.clone(),
837 session_event_senders.clone(),
838 Some(account_sink.inbox()),
839 Some(data_dir.clone()),
840 config.clone(),
841 provider_registry.clone(),
842 permission_checker.clone(),
843 project_store.clone(),
844 )
845 .await
846 .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?,
847 );
848
849 let child_adapter = Arc::new(crate::tools::ChildSessionAdapter {
856 session_store: session_store.clone(),
857 storage: storage.clone(),
858 persistence: persistence.clone(),
859 session_messenger: Some(session_messenger.clone()),
860 scheduler: spawn_scheduler.clone(),
861 sessions_cache: sessions.clone(),
862 agent_runners: agent_runners.clone(),
863 session_event_senders: session_event_senders.clone(),
864 subagent_model_resolver: None,
865 config: config.clone(),
866 project_store: Some(project_store.clone()),
867 parent_wait_slots: Arc::new(dashmap::DashMap::new()),
868 });
869 let guardian_spawner: Arc<dyn bamboo_engine::GuardianSpawner> = child_adapter.clone();
870 child_completion_coordinator
873 .set_guardian_spawner(guardian_spawner.clone())
874 .await;
875
876 let bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook> =
880 child_completion_coordinator.clone();
881
882 child_completion_coordinator.spawn_child_wait_watchdog();
888
889 reconcile_fabric_on_boot(&config, &bamboo_home_dir).await;
894
895 let service_manager = Arc::new(crate::service_manager::ServiceManager::new());
901 let boot_reconcile_services_handle = {
913 let service_manager = service_manager.clone();
914 let app_data_dir = bamboo_home_dir.clone();
915 tokio::spawn(async move {
916 crate::plugin_installer::boot_reconcile_services(&app_data_dir, &service_manager)
917 .await;
918 })
919 };
920
921 let (config_watcher, config_live_health, mcp_config_live_health) =
922 super::config_runtime::ConfigWatcherRuntime::start(
923 bamboo_home_dir.clone(),
924 config.clone(),
925 config_facade.clone(),
926 config_io_lock.clone(),
927 provider_registry.clone(),
928 provider_lock.clone(),
929 mcp_manager.clone(),
930 account_sink.clone(),
931 );
932 let credential_store = Arc::new(bamboo_config::CredentialStore::open(&bamboo_home_dir));
933
934 Ok(Self {
935 app_data_dir: bamboo_home_dir,
936 config,
937 config_facade,
938 config_io_lock,
939 config_live_health,
940 mcp_config_live_health,
941 config_watcher,
942 project_resource_watcher,
943 credential_store,
944 fabric_deployer,
945 embedded_broker,
946 health_monitor,
947 provider: provider_lock,
948 provider_handle,
949 sessions,
950 storage,
951 session_store,
952 project_store,
953 project_context_resolver,
954 session_repo,
955 persistence,
956 session_inbox,
957 session_activation_router,
958 session_messenger,
959 spawn_scheduler,
960 child_completion_coordinator,
961 guardian_spawner,
962 bash_resume_hook,
963 schedule_store,
964 schedule_manager,
965 connect_manager,
966 tool_factory,
967 permission_checker,
968 permission_section,
969 permission_io_lock,
970 approval_registry,
971 notification_service,
972 session_watchers,
973 cancel_tokens: Arc::new(RwLock::new(HashMap::new())),
974 mcp_proxy_shutdown,
975 skill_manager,
976 workflow_runs,
977 mcp_manager,
978 service_manager,
979 boot_reconcile_services_handle: tokio::sync::Mutex::new(Some(
980 boot_reconcile_services_handle,
981 )),
982 metrics_service,
983 agent_runners,
984 execute_startups: Arc::new(std::sync::Mutex::new(HashMap::new())),
985 session_event_senders,
986 account_sink,
987 process_registry,
988 metrics_bus: None, agent,
990 provider_registry,
991 provider_router,
992 model_catalog,
993 title_gen_in_flight: Arc::new(dashmap::DashSet::new()),
994 pairing_codes: Arc::new(dashmap::DashMap::new()),
995 pairing_code_guard: Arc::new(crate::handlers::settings::PairingCodeGuard::default()),
996 root_password_guard: Arc::new(crate::handlers::settings::RootPasswordGuard::default()),
997 codex_run_tokens,
998 })
1000 }
1001}
1002
1003pub struct EmbeddedBroker {
1006 task: tokio::task::JoinHandle<()>,
1007 gc_task: tokio::task::JoinHandle<()>,
1008}
1009
1010impl Drop for EmbeddedBroker {
1011 fn drop(&mut self) {
1012 self.task.abort();
1013 self.gc_task.abort();
1014 }
1015}
1016
1017pub struct HealthMonitor(tokio::task::JoinHandle<()>);
1021
1022impl Drop for HealthMonitor {
1023 fn drop(&mut self) {
1024 self.0.abort();
1025 }
1026}
1027
1028async fn maybe_embed_broker(
1040 config: &mut bamboo_llm::Config,
1041 data_dir: &std::path::Path,
1042) -> Option<EmbeddedBroker> {
1043 if let Some(external) = load_external_broker(data_dir) {
1048 let endpoint = external.endpoint.trim().to_string();
1049 if broker_endpoint_reachable(&endpoint).await {
1050 tracing::info!(%endpoint, "using external broker from broker.json");
1051 config.subagents_mut().broker = Some(external);
1052 return None;
1053 }
1054 tracing::warn!(
1055 %endpoint,
1056 "broker.json endpoint is unreachable — embedding a fresh in-process broker instead"
1057 );
1058 }
1059
1060 let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
1061 Ok(l) => l,
1062 Err(e) => {
1063 tracing::warn!("embedded broker: bind failed, sub-agent dispatch disabled: {e}");
1064 return None;
1065 }
1066 };
1067 let port = match listener.local_addr() {
1068 Ok(a) => a.port(),
1069 Err(e) => {
1070 tracing::warn!("embedded broker: local_addr failed: {e}");
1071 return None;
1072 }
1073 };
1074 let token = uuid::Uuid::new_v4().simple().to_string();
1075 let root = data_dir.join("broker");
1076 let core = Arc::new(bamboo_broker::BrokerCore::new(root));
1077 let gc_task = core
1080 .clone()
1081 .spawn_mailbox_gc(std::time::Duration::from_secs(300));
1082 let server = Arc::new(bamboo_broker::BrokerServer::new(core, token.clone()));
1083
1084 let task = tokio::spawn(async move {
1085 if let Err(e) = server.serve(listener).await {
1086 tracing::error!("embedded broker serve loop ended: {e}");
1087 }
1088 });
1089
1090 config.subagents_mut().broker = Some(bamboo_config::BrokerClientConfig {
1093 endpoint: format!("ws://127.0.0.1:{port}"),
1094 token,
1095 token_encrypted: None,
1096 credential_ref: None,
1097 configured: false,
1098 });
1099 tracing::info!(port, "embedded mailbox bus (broker) started in-process");
1100 Some(EmbeddedBroker { task, gc_task })
1101}
1102
1103fn load_external_broker(data_dir: &std::path::Path) -> Option<bamboo_config::BrokerClientConfig> {
1112 let path = data_dir.join("broker.json");
1113 if let Err(error) = bamboo_config::migrate_external_broker_credentials(data_dir)
1114 .and_then(|_| bamboo_config::ensure_provider_mcp_migration_ready(data_dir))
1115 {
1116 tracing::warn!(error = %error, "external broker credential migration unavailable");
1117 return None;
1118 }
1119 let bytes = std::fs::read(&path).ok()?;
1120 match parse_external_broker_snapshot(&bytes, data_dir) {
1121 Ok(cfg) => Some(cfg),
1122 Err(error) => {
1123 tracing::warn!(?path, %error, "broker.json is unavailable");
1124 None
1125 }
1126 }
1127}
1128
1129fn parse_external_broker_snapshot(
1130 bytes: &[u8],
1131 data_dir: &std::path::Path,
1132) -> bamboo_config::ConfigStoreResult<bamboo_config::BrokerClientConfig> {
1133 let mut cfg: bamboo_config::BrokerClientConfig = serde_json::from_slice(bytes)?;
1134 if !cfg.token.trim().is_empty()
1135 || cfg
1136 .token_encrypted
1137 .as_deref()
1138 .is_some_and(|value| !value.trim().is_empty())
1139 {
1140 return Err(bamboo_config::ConfigStoreError::Validation(
1141 "legacy broker credential appeared after migration".to_string(),
1142 ));
1143 }
1144 if cfg.endpoint.trim().is_empty() {
1145 return Err(bamboo_config::ConfigStoreError::Validation(
1146 "broker endpoint is empty".to_string(),
1147 ));
1148 }
1149 cfg.hydrate_credential_from_store(data_dir)?;
1150 Ok(cfg)
1151}
1152
1153async fn broker_endpoint_reachable(endpoint: &str) -> bool {
1158 let host_port = endpoint
1159 .trim()
1160 .trim_start_matches("wss://")
1161 .trim_start_matches("ws://")
1162 .split('/')
1163 .next()
1164 .unwrap_or("");
1165 if host_port.is_empty() {
1166 return false;
1167 }
1168 matches!(
1169 tokio::time::timeout(
1170 std::time::Duration::from_millis(500),
1171 tokio::net::TcpStream::connect(host_port),
1172 )
1173 .await,
1174 Ok(Ok(_))
1175 )
1176}
1177
1178async fn reconcile_fabric_on_boot(
1183 config: &Arc<RwLock<bamboo_llm::Config>>,
1184 data_dir: &std::path::Path,
1185) {
1186 use bamboo_config::cluster_fabric::NodeStatus;
1187
1188 let snapshot = {
1189 let mut cfg = config.write().await;
1190 let mut changed = 0usize;
1191 for node in &mut cfg.cluster_fabric.nodes {
1192 if let Some(state) = node.state.as_mut() {
1193 if matches!(state.status, NodeStatus::Running | NodeStatus::Deploying) {
1194 state.status = NodeStatus::Unreachable;
1195 state.last_error =
1196 Some("orchestrator restarted; worker no longer tracked".to_string());
1197 changed += 1;
1198 }
1199 }
1200 }
1201 if changed == 0 {
1202 return;
1203 }
1204 tracing::info!(
1205 reconciled = changed,
1206 "cluster-fabric: marked stale Running nodes Unreachable on boot"
1207 );
1208 cfg.clone()
1209 };
1210
1211 if let Err(e) = snapshot.save_to_dir(data_dir.to_path_buf()) {
1212 tracing::warn!("cluster-fabric boot reconcile: failed to persist: {e}");
1213 }
1214}
1215
1216#[cfg(test)]
1217mod broker_embed_tests {
1218 use super::{broker_endpoint_reachable, load_external_broker, parse_external_broker_snapshot};
1219
1220 #[test]
1221 fn load_external_broker_reads_broker_json_not_config() {
1222 let _key = bamboo_config::encryption::set_test_encryption_key([0x57; 32]);
1223 let dir = tempfile::tempdir().unwrap();
1224 assert!(load_external_broker(dir.path()).is_none());
1226
1227 std::fs::write(
1229 dir.path().join("broker.json"),
1230 r#"{ "endpoint": "wss://broker.example:9600", "token": "t" }"#,
1231 )
1232 .unwrap();
1233 let got = load_external_broker(dir.path()).expect("parsed");
1234 assert_eq!(got.endpoint, "wss://broker.example:9600");
1235 assert_eq!(got.token, "t");
1236 assert_eq!(
1237 got.credential_ref.as_ref().unwrap().as_str(),
1238 "broker.external.bearer_token"
1239 );
1240 assert!(got.configured);
1241 let durable = std::fs::read_to_string(dir.path().join("broker.json")).unwrap();
1242 assert!(!durable.contains("\"token\""));
1243 assert!(!durable.contains("token_encrypted"));
1244 assert!(!durable.contains("\"t\""));
1245
1246 std::fs::write(dir.path().join("broker.json"), r#"{ "endpoint": " " }"#).unwrap();
1248 assert!(load_external_broker(dir.path()).is_none());
1249
1250 std::fs::write(dir.path().join("broker.json"), "not json").unwrap();
1252 assert!(load_external_broker(dir.path()).is_none());
1253 }
1254
1255 #[test]
1256 fn external_broker_reference_fails_closed_and_tracks_generic_cas_updates() {
1257 let _key = bamboo_config::encryption::set_test_encryption_key([0x58; 32]);
1258 let dir = tempfile::tempdir().unwrap();
1259 let reference =
1260 bamboo_config::credential_ref("broker", "external", "bearer_token").unwrap();
1261 std::fs::write(
1262 dir.path().join("broker.json"),
1263 serde_json::to_vec_pretty(&serde_json::json!({
1264 "endpoint": "wss://broker.example:9600",
1265 "credential_ref": reference,
1266 "configured": true,
1267 }))
1268 .unwrap(),
1269 )
1270 .unwrap();
1271
1272 assert!(load_external_broker(dir.path()).is_none());
1273 std::fs::write(
1274 dir.path().join("broker.json"),
1275 serde_json::to_vec_pretty(&serde_json::json!({
1276 "endpoint": "wss://broker.example:9600",
1277 "credential_ref": reference,
1278 "configured": false,
1279 }))
1280 .unwrap(),
1281 )
1282 .unwrap();
1283 assert!(load_external_broker(dir.path()).is_none());
1284 let store = bamboo_config::CredentialStore::open(dir.path());
1285 store
1286 .replace(
1287 reference.clone(),
1288 "replacement-token",
1289 bamboo_config::CredentialSource::User,
1290 0,
1291 )
1292 .unwrap();
1293 let hydrated = load_external_broker(dir.path()).unwrap();
1294 assert_eq!(hydrated.token, "replacement-token");
1295 store.clear(&reference, 1).unwrap();
1296 assert!(load_external_broker(dir.path()).is_none());
1297 }
1298
1299 #[test]
1300 fn broker_snapshot_with_late_legacy_secret_fails_closed() {
1301 let _key = bamboo_config::encryption::set_test_encryption_key([0x59; 32]);
1302 let dir = tempfile::tempdir().unwrap();
1303 let error = parse_external_broker_snapshot(
1304 br#"{
1305 "endpoint": "wss://broker.example:9600",
1306 "token": "late-legacy-token"
1307 }"#,
1308 dir.path(),
1309 )
1310 .unwrap_err();
1311 assert!(error.to_string().contains("appeared after migration"));
1312 }
1313
1314 #[tokio::test]
1315 async fn reachability_probe_distinguishes_live_from_dead() {
1316 let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1319 let dead = l.local_addr().unwrap();
1320 drop(l);
1321 assert!(!broker_endpoint_reachable(&format!("ws://{dead}")).await);
1322
1323 assert!(!broker_endpoint_reachable("").await);
1325 assert!(!broker_endpoint_reachable("ws://").await);
1326
1327 let live = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1329 let addr = live.local_addr().unwrap();
1330 assert!(broker_endpoint_reachable(&format!("ws://{addr}/stream")).await);
1331 }
1332}