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