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::tool_event_router::{CombinedToolEventPublisher, ToolEventRouter};
9use crate::tools::OptionalSubagentModelResolver;
10use bamboo_agent_core::storage::Storage;
11#[cfg(not(test))]
12use bamboo_memory::memory_store::{resolve_jiandu_data_root, BAMBOO_JIANDU_DATA_DIR_ENV};
13use bamboo_plugin_protocol::{NoopToolEventPublisher, ToolEventPublisher};
14
15fn default_app_state_memory_store(
16 bamboo_home_dir: &std::path::Path,
17) -> Result<bamboo_memory::memory_store::MemoryStore, AppError> {
18 #[cfg(test)]
19 {
20 let root = bamboo_home_dir.join("jiandu");
21 tracing::info!(
22 target: "bamboo.memory",
23 mode = "test-isolated",
24 root = %root.display(),
25 "selected Jiandu data root"
26 );
27 Ok(bamboo_memory::memory_store::MemoryStore::new(root))
28 }
29 #[cfg(not(test))]
30 {
31 let _ = bamboo_home_dir;
32 let selection = resolve_jiandu_data_root(std::env::var_os(BAMBOO_JIANDU_DATA_DIR_ENV))
33 .map_err(|message| AppError::InternalError(anyhow::anyhow!(message)))?;
34 let mode = selection.mode();
35 let root = selection.into_path();
36 tracing::info!(
37 target: "bamboo.memory",
38 mode,
39 root = %root.display(),
40 "selected Jiandu data root"
41 );
42 Ok(bamboo_memory::memory_store::MemoryStore::new(root))
43 }
44}
45
46impl AppState {
47 pub async fn new(bamboo_home_dir: PathBuf) -> Result<Self, AppError> {
89 let memory_store = default_app_state_memory_store(&bamboo_home_dir)?;
90 Self::new_with_memory_store(bamboo_home_dir, memory_store).await
91 }
92
93 pub async fn new_with_memory_store(
96 bamboo_home_dir: PathBuf,
97 memory_store: bamboo_memory::memory_store::MemoryStore,
98 ) -> Result<Self, AppError> {
99 bamboo_config::paths::init_bamboo_dir(bamboo_home_dir.clone());
102
103 let (config, config_facade) = match bamboo_config::ConfigFacade::open_or_migrate(
110 &bamboo_home_dir,
111 ) {
112 Ok(facade) => {
113 let facade = Arc::new(facade);
114 let config =
115 super::config_runtime::load_facade_effective_config(&facade, &bamboo_home_dir);
116 (config, Some(facade))
117 }
118 Err(error) => {
119 if bamboo_config::modular_authority_boundary_present(&bamboo_home_dir)
120 .unwrap_or(true)
121 {
122 return Err(AppError::InternalError(anyhow::anyhow!(
123 "modular configuration authority is unavailable: {error}"
124 )));
125 }
126 tracing::warn!(
127 error = %error,
128 "modular configuration facade is unavailable; retaining recovered legacy authority"
129 );
130 (
131 Config::from_data_dir_without_publish(Some(bamboo_home_dir.clone())),
132 None,
133 )
134 }
135 };
136 config.publish_env_vars();
137
138 if config.plugin_trust.enforcement_is_off() {
150 super::config_runtime::warn_plugin_trust_enforcement_off();
151 }
152
153 let provider_registry =
154 match bamboo_llm::ProviderRegistry::from_config(&config, bamboo_home_dir.clone()).await
155 {
156 Ok(registry) => Arc::new(registry),
157 Err(e) => {
158 tracing::error!("Failed to create provider registry: {}", e);
159 Arc::new(
160 bamboo_llm::ProviderRegistry::from_config(
161 &Config::default(),
162 bamboo_home_dir.clone(),
163 )
164 .await
165 .expect("Cannot create even an empty provider registry"),
166 )
167 }
168 };
169
170 let provider = provider_registry.get_default().unwrap_or_else(|| {
171 let default_provider_name = provider_registry.default_provider_name();
172 let message = if config.has_provider_instances() {
173 format!(
174 "Default provider instance '{}' is not available or failed to initialize",
175 default_provider_name
176 )
177 } else {
178 format!(
179 "Provider '{}' is not available or failed to initialize",
180 config.provider
181 )
182 };
183 Arc::new(UnconfiguredProvider { message }) as Arc<dyn LLMProvider>
184 });
185
186 Self::new_with_provider_and_facade(
187 bamboo_home_dir,
188 config,
189 provider,
190 config_facade,
191 Arc::new(NoopToolEventPublisher),
192 memory_store,
193 )
194 .await
195 }
196
197 pub async fn new_with_provider(
212 bamboo_home_dir: PathBuf,
213 config: Config,
214 provider: Arc<dyn LLMProvider>,
215 ) -> Result<Self, AppError> {
216 let memory_store = default_app_state_memory_store(&bamboo_home_dir)?;
217 Self::new_with_provider_and_facade(
218 bamboo_home_dir,
219 config,
220 provider,
221 None,
222 Arc::new(NoopToolEventPublisher),
223 memory_store,
224 )
225 .await
226 }
227
228 pub async fn new_with_provider_and_tool_event_publisher(
231 bamboo_home_dir: PathBuf,
232 config: Config,
233 provider: Arc<dyn LLMProvider>,
234 tool_event_publisher: Arc<dyn ToolEventPublisher>,
235 ) -> Result<Self, AppError> {
236 let memory_store = default_app_state_memory_store(&bamboo_home_dir)?;
237 Self::new_with_provider_and_facade(
238 bamboo_home_dir,
239 config,
240 provider,
241 None,
242 tool_event_publisher,
243 memory_store,
244 )
245 .await
246 }
247
248 async fn new_with_provider_and_facade(
249 bamboo_home_dir: PathBuf,
250 config: Config,
251 provider: Arc<dyn LLMProvider>,
252 config_facade: Option<Arc<bamboo_config::ConfigFacade>>,
253 tool_event_publisher: Arc<dyn ToolEventPublisher>,
254 memory_store: bamboo_memory::memory_store::MemoryStore,
255 ) -> Result<Self, AppError> {
256 let data_dir = bamboo_home_dir.clone();
258 let (session_store, storage) = init_storage(&data_dir).await?;
259 let session_create_operations =
260 Arc::new(super::session_create_operations::SessionCreateOperationStore::new(&data_dir));
261 let mutation_idempotency =
262 Arc::new(super::mutation_idempotency::MutationIdempotencyStore::default());
263 match session_create_operations.prune_expired().await {
264 Ok(0) => {}
265 Ok(deleted) => tracing::info!(
266 target: "bamboo.session_create",
267 phase = "retention_cleanup",
268 outcome = "expired_pruned",
269 deleted,
270 "pruned expired session-create operation receipts"
271 ),
272 Err(error) => tracing::warn!(
273 target: "bamboo.session_create",
274 phase = "retention_cleanup",
275 outcome = "cleanup_failed",
276 error = %error,
277 "failed to prune expired session-create operation receipts"
278 ),
279 }
280 let project_store = Arc::new(bamboo_projects::ProjectStore::open(&data_dir).map_err(
281 |error| {
282 AppError::InternalError(anyhow::anyhow!(
283 "failed to initialize Project registry: {error}"
284 ))
285 },
286 )?);
287 let persistence = Arc::new(LockedSessionStore::new(storage.clone()));
288 let session_inbox: Arc<dyn bamboo_domain::SessionInboxPort> =
289 Arc::new(bamboo_storage::FileSessionInbox::new(
290 session_store.clone(),
291 bamboo_domain::SessionInboxLimits::default(),
292 ));
293 let session_activation_router = bamboo_engine::SessionActivationRouter::new();
294 let session_messenger = Arc::new(bamboo_engine::SessionMessenger::new(
295 storage.clone(),
296 session_inbox.clone(),
297 session_activation_router.clone(),
298 ));
299
300 let sessions: bamboo_engine::SessionCache = Arc::default();
302
303 let mut config = config;
310 let embedded_broker = maybe_embed_broker(&mut config, &data_dir).await;
311
312 let config = Arc::new(RwLock::new(config));
313
314 let live_default_workspace: Arc<dyn Fn() -> Option<PathBuf> + Send + Sync> = {
325 let config_for_workspace = config.clone();
326 let last_known: Arc<std::sync::Mutex<Option<PathBuf>>> =
327 Arc::new(std::sync::Mutex::new(None));
328 Arc::new(move || match config_for_workspace.try_read() {
329 Ok(cfg) => {
330 let path = cfg.get_default_work_area_path();
331 if let Ok(mut cache) = last_known.lock() {
332 *cache = path.clone();
333 }
334 path
335 }
336 Err(_) => last_known.lock().ok().and_then(|cache| cache.clone()),
337 })
338 };
339
340 let live_workspace_root: Arc<
345 dyn Fn() -> bamboo_agent_core::workspace_state::WorkspaceRootConfig + Send + Sync,
346 > = {
347 let app_data_dir = data_dir.clone();
348 Arc::new(
349 move || bamboo_agent_core::workspace_state::WorkspaceRootConfig {
350 root: bamboo_config::paths::resolve_workspace_root_in(&app_data_dir),
351 confine: bamboo_config::paths::workspace_confinement_enforced(),
352 },
353 )
354 };
355
356 let workspace_resolver = bamboo_agent_core::workspace_state::WorkspaceResolver::new(
357 {
358 let provider = live_default_workspace.clone();
359 move || provider()
360 },
361 {
362 let provider = live_workspace_root.clone();
363 move || provider()
364 },
365 );
366
367 bamboo_agent_core::workspace_state::set_default_workspace_provider(Box::new({
368 let provider = live_default_workspace;
369 move || provider()
370 }));
371 bamboo_agent_core::workspace_state::set_workspace_root_provider(Box::new({
372 let provider = live_workspace_root;
373 move || provider()
374 }));
375
376 let (permission_checker, permission_section) =
377 load_permission_checker(&bamboo_home_dir).await?;
378 let permission_io_lock = Arc::new(tokio::sync::Mutex::new(()));
379 let notification_service = Arc::new(bamboo_notification::NotificationService::new(
380 bamboo_home_dir.join("notification_preferences.json"),
381 ));
382 let session_watchers = super::watchers::SessionWatchers::new();
383 let (mcp_manager, _legacy_mcp_bootstrap) =
384 init_mcp_manager(config.clone(), &bamboo_home_dir);
385 let skill_manager = init_skill_manager(&data_dir).await;
386 let metrics_service = init_metrics_service(&data_dir).await?;
387
388 let service_manager = Arc::new(crate::service_manager::ServiceManager::new());
393 let tool_event_router = ToolEventRouter::new(service_manager.clone());
394 let tool_event_publisher: Arc<dyn ToolEventPublisher> = Arc::new(
395 CombinedToolEventPublisher::new(tool_event_router.clone(), tool_event_publisher),
396 );
397
398 let startup_sessions = {
399 let entries = session_store.list_index_entries().await;
400 let mut sessions = Vec::new();
401 for entry in entries {
402 if let Some(session) = session_store
403 .load_session(&entry.id)
404 .await
405 .map_err(AppError::StorageError)?
406 {
407 sessions.push(session);
408 }
409 }
410 sessions
411 };
412 metrics_service
413 .reconcile_startup_sessions(startup_sessions, &[])
414 .await
415 .map_err(|error| {
416 AppError::InternalError(anyhow::anyhow!(
417 "Failed to reconcile stale metrics state on startup: {error}"
418 ))
419 })?;
420
421 let agent_runners: Arc<RwLock<HashMap<String, AgentRunner>>> =
422 Arc::new(RwLock::new(HashMap::new()));
423 let process_registry = Arc::new(ProcessRegistry::new());
428 let (provider_lock, provider_handle) = build_provider_handles(provider);
429
430 let config_snapshot = config.read().await;
432 let provider_registry = match bamboo_llm::ProviderRegistry::from_config(
433 &config_snapshot,
434 bamboo_home_dir.clone(),
435 )
436 .await
437 {
438 Ok(registry) => Arc::new(registry),
439 Err(e) => {
440 tracing::error!("Failed to create provider registry: {}", e);
441 Arc::new(
442 bamboo_llm::ProviderRegistry::from_config(
443 &Config::default(),
444 bamboo_home_dir.clone(),
445 )
446 .await
447 .expect("Cannot create even an empty provider registry"),
448 )
449 }
450 };
451 drop(config_snapshot);
452
453 let provider_router = Arc::new(bamboo_llm::ProviderModelRouter::new(
454 provider_registry.clone(),
455 ));
456 let model_catalog = Arc::new(bamboo_llm::ModelCatalogService::new(
457 provider_registry.clone(),
458 ));
459
460 let session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>> =
465 Arc::new(RwLock::new(HashMap::new()));
466
467 let notification_relay_deps = crate::app_state::session_events::NotificationRelayDeps {
473 notification_service: notification_service.clone(),
474 session_event_senders: session_event_senders.clone(),
475 session_watchers: session_watchers.clone(),
476 config: config.clone(),
477 };
478
479 let ledger_schedule_bridge =
482 Arc::new(crate::schedule_app::LateBoundLedgerBridge::default());
483
484 let session_repo = bamboo_engine::SessionRepository::new(
488 sessions.clone(),
489 storage.clone(),
490 persistence.clone(),
491 );
492
493 let account_sink = bamboo_engine::events::AccountEventSink::new(data_dir.join("events"))
497 .map_err(|e| {
498 AppError::InternalError(anyhow::anyhow!(
499 "failed to initialize account change-feed journal: {e}"
500 ))
501 })?;
502 let project_resource_watcher = super::project_watcher::ProjectResourceWatcher::start(
503 project_store.clone(),
504 account_sink.clone(),
505 std::time::Duration::from_millis(120),
506 )
507 .map_err(|error| {
508 AppError::InternalError(anyhow::anyhow!(
509 "failed to start Project resource watcher: {error}"
510 ))
511 })?;
512
513 let base_tools = build_base_tools(
514 config.clone(),
515 permission_checker.clone(),
516 mcp_manager.clone(),
517 skill_manager.clone(),
518 session_repo.clone(),
519 session_store.clone(),
520 storage.clone(),
521 bamboo_home_dir.clone(),
522 notification_service.clone(),
523 session_event_senders.clone(),
524 session_watchers.clone(),
525 ledger_schedule_bridge.clone(),
526 project_store.clone(),
527 account_sink.clone(),
528 workspace_resolver.clone(),
529 tool_event_publisher.clone(),
530 memory_store.clone(),
531 );
532
533 let workflow_runs = crate::workflow::WorkflowRunAccess::new_with_permission_config(
537 &data_dir,
538 base_tools.clone(),
539 skill_manager.clone(),
540 session_repo.clone(),
541 permission_checker.permission_config(),
542 )
543 .await
544 .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?;
545
546 spawn_session_map_cleanup_task(
550 agent_runners.clone(),
551 session_event_senders.clone(),
552 session_watchers.clone(),
553 sessions.clone(),
554 None,
555 );
556
557 {
563 let mut workflow_events = skill_manager.store().subscribe_workflow_catalog();
564 let account_sink = account_sink.clone();
565 tokio::spawn(async move {
566 loop {
567 let event = match workflow_events.recv().await {
568 Ok(event) => event,
569 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
570 tracing::warn!("Workflow catalog event bridge lagged by {skipped}");
571 continue;
572 }
573 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
574 };
575 let event = match event.kind {
576 bamboo_skills::WorkflowCatalogEventKind::Changed => {
577 AgentEvent::WorkflowChanged {
578 workflow_id: event.workflow_id,
579 revision: event.revision,
580 scope: event.scope,
581 }
582 }
583 bamboo_skills::WorkflowCatalogEventKind::Invalid => {
584 AgentEvent::WorkflowInvalid {
585 workflow_id: event.workflow_id,
586 revision: event.revision,
587 scope: event.scope,
588 }
589 }
590 bamboo_skills::WorkflowCatalogEventKind::Recovered => {
591 AgentEvent::WorkflowRecovered {
592 workflow_id: event.workflow_id,
593 revision: event.revision,
594 scope: event.scope,
595 }
596 }
597 };
598 account_sink.record(None, &event);
599 }
600 });
601 }
602 let (approval_registry, restart_approval_events) =
603 bamboo_engine::external_agents::live::initialize_durable_approvals(
604 data_dir.join("approvals/child-approvals-v1.json"),
605 )
606 .map_err(|error| {
607 AppError::InternalError(anyhow::anyhow!(
608 "failed to initialize durable child approvals: {error}"
609 ))
610 })?;
611 for event in restart_approval_events {
612 account_sink.record(event.session_id(), &event);
613 }
614
615 let child_tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = base_tools.clone();
618
619 let project_context_resolver = Arc::new(
624 bamboo_engine::project_context::ProjectContextResolver::new_with_workspace_resolver(
625 Arc::new(crate::project_context::ProjectStoreContextSource::new(
626 project_store.clone(),
627 )),
628 workspace_resolver.clone(),
629 ),
630 );
631 let mut agent_builder = bamboo_engine::Agent::builder()
632 .storage(storage.clone())
633 .persistence(Arc::new(session_repo.clone()))
634 .session_inbox(session_inbox.clone())
635 .activation_router(session_activation_router.clone())
636 .session_messenger(session_messenger.clone())
637 .attachment_reader(session_store.clone())
638 .skill_manager(skill_manager.clone())
639 .metrics_collector(metrics_service.collector())
640 .config(config.clone())
641 .provider(provider_handle.clone())
642 .memory_store(memory_store.clone())
643 .default_tools(base_tools.clone())
644 .project_context_resolver(project_context_resolver.clone());
645 if let Some(permission_config) = permission_checker.permission_config() {
646 agent_builder = agent_builder.permission_config(permission_config);
647 }
648 let agent = Arc::new(
649 agent_builder
650 .build()
651 .expect("agent runtime should be fully configured"),
652 );
653
654 let child_completion_coordinator =
655 Arc::new(bamboo_engine::ChildCompletionCoordinator::new(
656 storage.clone(),
657 persistence.clone(),
658 sessions.clone(),
659 agent_runners.clone(),
660 session_event_senders.clone(),
661 agent.clone(),
662 config.clone(),
663 provider_registry.clone(),
664 provider_router.clone(),
665 data_dir.clone(),
666 Some(account_sink.inbox()),
667 ));
668 session_activation_router
669 .set_spawner(child_completion_coordinator.clone())
670 .await;
671
672 let config_snapshot = config.read().await.clone();
674
675 let mcp_proxy_shutdown = tokio_util::sync::CancellationToken::new();
682 if let Some(broker) = config_snapshot.subagents().broker.clone() {
683 if !broker.endpoint.trim().is_empty() {
684 let backend: std::sync::Arc<dyn bamboo_agent_core::tools::ToolExecutor> =
685 std::sync::Arc::new(bamboo_mcp::executor::McpToolExecutor::new(
686 mcp_manager.clone(),
687 mcp_manager.tool_index(),
688 ));
689 let shutdown = mcp_proxy_shutdown.clone();
690
691 let role_entries: Vec<(String, Vec<String>)> = config_snapshot
700 .subagents()
701 .mcp_role_allowlist
702 .iter()
703 .map(|e| (e.role.clone(), e.tools.clone()))
704 .collect();
705 if role_entries.is_empty() {
706 tracing::info!(
712 "mcp proxy: no subagents.mcp_role_allowlist configured — every worker \
713 role sees/can call the full host-bound MCP tool set (opt in a role \
714 policy in config.json to scope tools per role; see issue #54)"
715 );
716 }
717 let known_tools: std::collections::HashSet<String> = backend
725 .list_tools()
726 .into_iter()
727 .map(|t| t.function.name)
728 .collect();
729 if !role_entries.is_empty() && known_tools.is_empty() {
730 tracing::warn!(
731 "mcp role allowlist: the orchestrator's MCP tool set was empty at \
732 policy-load time (servers may still be connecting in the background) — \
733 skipped tool-name typo validation for subagents.mcp_role_allowlist"
734 );
735 }
736 let allowlist = std::sync::Arc::new(bamboo_broker::RoleToolAllowlist::from_config(
737 role_entries,
738 &known_tools,
739 ));
740 tokio::spawn(async move {
741 let me = bamboo_broker::AgentRef {
742 session_id: bamboo_broker::ORCHESTRATOR_ID.to_string(),
743 role: Some("orchestrator".into()),
744 };
745 bamboo_broker::serve_mcp_proxy_supervised(
746 &broker.endpoint,
747 me,
748 &broker.token,
749 backend,
750 allowlist,
751 shutdown,
752 )
753 .await;
754 });
755 }
756 }
757 let parent_approval_reviewer = Arc::new(
758 crate::app_state::parent_approval_reviewer::ParentAgentApprovalReviewer::new(
759 session_repo.clone(),
760 provider_router.clone(),
761 ),
762 );
763 let codex_run_tokens = Arc::new(crate::codex_run_tokens::CodexRunTokenRegistry::default());
764 let external_runner =
765 bamboo_engine::external_agents::runtime::build_external_child_runner_with_codex_tokens(
766 &config_snapshot,
767 Some(approval_registry.clone()),
768 Some(parent_approval_reviewer),
769 permission_checker.permission_config(),
770 Some(codex_run_tokens.clone()),
771 );
772 external_runner.set_session_inbox_runtime(Some(
773 bamboo_engine::execution::spawn::SessionInboxRuntimeBinding {
774 router: session_activation_router.clone(),
775 inbox: session_inbox.clone(),
776 storage: storage.clone(),
777 persistence: persistence.clone(),
778 },
779 ));
780 let spawn_scheduler = build_spawn_scheduler(
781 agent.clone(),
782 child_tools,
783 sessions.clone(),
784 agent_runners.clone(),
785 session_event_senders.clone(),
786 external_runner,
787 Some(provider_router.clone()),
788 Some(child_completion_coordinator.clone()),
789 Some(data_dir.clone()),
790 Some(account_sink.inbox()),
791 Some(Arc::new(
792 crate::app_state::session_events::NotificationRelayLaunchHook::new(
793 notification_relay_deps.clone(),
794 ),
795 )),
796 );
797 child_completion_coordinator
798 .set_spawn_scheduler(&spawn_scheduler)
799 .await;
800
801 let tools_with_task = base_tools.clone();
802
803 let schedule_store = init_schedule_store(&data_dir).await?;
804
805 ledger_schedule_bridge
807 .bind(Arc::new(crate::schedule_app::ScheduleLedgerBridge::new(
808 schedule_store.clone(),
809 )))
810 .await;
811
812 let schedule_manager = build_schedule_manager(
813 schedule_store.clone(),
814 agent.clone(),
815 tools_with_task.clone(),
816 permission_checker.permission_config(),
817 sessions.clone(),
818 agent_runners.clone(),
819 session_event_senders.clone(),
820 persistence.clone(),
821 config.clone(),
822 provider_registry.clone(),
823 Some(data_dir.clone()),
824 Some(account_sink.inbox()),
825 notification_relay_deps.clone(),
826 project_store.clone(),
827 workspace_resolver.clone(),
828 );
829
830 bamboo_engine::auto_dream::spawn_auto_dream_task_with_project_resolver(
831 bamboo_engine::auto_dream::AutoDreamContext {
832 session_store: session_store.clone(),
833 storage: storage.clone(),
834 memory: memory_store.clone(),
835 provider: provider_handle.clone(),
836 config: config.clone(),
837 provider_registry: provider_registry.clone(),
838 },
839 project_context_resolver.as_ref().clone(),
840 );
841
842 bamboo_engine::gardener::spawn_gardener_task_with_project_resolver(
846 bamboo_engine::auto_dream::AutoDreamContext {
847 session_store: session_store.clone(),
848 storage: storage.clone(),
849 memory: memory_store.clone(),
850 provider: provider_handle.clone(),
851 config: config.clone(),
852 provider_registry: provider_registry.clone(),
853 },
854 project_context_resolver.clone(),
855 );
856
857 bamboo_engine::ledger_gardener::spawn_ledger_gardener_task(
861 bamboo_engine::ledger_gardener::LedgerGardenerContext {
862 dream: bamboo_engine::auto_dream::AutoDreamContext {
863 session_store: session_store.clone(),
864 storage: storage.clone(),
865 memory: memory_store.clone(),
866 provider: provider_handle.clone(),
867 config: config.clone(),
868 provider_registry: provider_registry.clone(),
869 },
870 schedule_bridge: Some(ledger_schedule_bridge.clone()),
871 },
872 );
873
874 let config_for_resolver = config.clone();
875 let subagent_model_resolver: OptionalSubagentModelResolver = {
876 let registry = provider_registry.clone();
877 Some(Arc::new(
878 move |subagent_type: String| -> futures::future::BoxFuture<
879 'static,
880 Option<bamboo_domain::ProviderModelRef>,
881 > {
882 let config_for_resolver = config_for_resolver.clone();
883 let registry = registry.clone();
884 Box::pin(async move {
885 let config_snap = config_for_resolver.read().await.clone();
886 bamboo_engine::model_config_helper::resolve_subagent_model_ref(
887 &config_snap,
888 &config_snap.provider,
889 ®istry,
890 &subagent_type,
891 )
892 })
893 },
894 ))
895 };
896
897 let config_io_lock = Arc::new(tokio::sync::Mutex::new(()));
901 let fabric_registry: crate::tools::DeployedRegistry =
902 Arc::new(tokio::sync::Mutex::new(HashMap::new()));
903 let fabric_bamboo_bin =
904 std::env::current_exe().unwrap_or_else(|_| std::path::PathBuf::from("bamboo"));
905 let credential_store = Arc::new(bamboo_config::CredentialStore::open(&bamboo_home_dir));
906 let mut fabric_deployer = bamboo_server_tools::FabricDeployer::new(
907 config.clone(),
908 config_io_lock.clone(),
909 bamboo_home_dir.clone(),
910 fabric_registry,
911 fabric_bamboo_bin,
912 );
913 if let Some(facade) = config_facade.clone() {
914 let event_sink = account_sink.clone();
915 fabric_deployer = fabric_deployer.with_modular_persistence(
916 facade,
917 credential_store.clone(),
918 Arc::new(move |event| {
919 let event_sink = event_sink.clone();
920 let event = event.clone();
921 Box::pin(async move {
922 super::config_runtime::publish_registry_event(&event_sink, &event).await;
923 })
924 }),
925 );
926 }
927 if let Err(error) = fabric_deployer.reconcile_stale_nodes_on_boot().await {
928 tracing::warn!(
929 error = %error,
930 "cluster-fabric boot reconcile failed; retaining the last adopted state"
931 );
932 }
933 let fabric_deployer = Arc::new(fabric_deployer);
934 let health_monitor = fabric_deployer
939 .clone()
940 .spawn_health_monitor()
941 .await
942 .map(HealthMonitor);
943
944 let tools = build_root_tools(
945 tools_with_task.clone(),
946 schedule_store.clone(),
947 schedule_manager.clone(),
948 session_store.clone(),
949 storage.clone(),
950 persistence.clone(),
951 session_messenger.clone(),
952 spawn_scheduler.clone(),
953 sessions.clone(),
954 agent_runners.clone(),
955 session_event_senders.clone(),
956 subagent_model_resolver,
957 config.clone(),
958 provider_registry.clone(),
959 config_snapshot.subagents().broker.clone(),
960 fabric_deployer.clone(),
961 project_store.clone(),
962 workspace_resolver.clone(),
963 );
964 let workflow_run_tool =
965 Arc::new(crate::workflow::WorkflowRunTool::new(workflow_runs.clone()));
966 let tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = Arc::new(
967 crate::tools::OverlayToolExecutor::new(tools, workflow_run_tool),
968 );
969
970 child_completion_coordinator
971 .set_root_tools(tools.clone())
972 .await;
973
974 for entry in session_store.list_index_entries().await {
979 match session_inbox.inspect(&entry.id).await {
980 Ok(backlog) if backlog.activation_pending() => {
981 if let Err(error) = bamboo_domain::SessionActivationPort::request_activation(
982 session_activation_router.as_ref(),
983 &entry.id,
984 backlog.activation_generation,
985 )
986 .await
987 {
988 tracing::error!(
989 session_id = %entry.id,
990 %error,
991 "failed to reactivate durable SessionInbox backlog during startup"
992 );
993 }
994 }
995 Ok(_) => {}
996 Err(error) => tracing::warn!(
997 session_id = %entry.id,
998 %error,
999 "failed to inspect SessionInbox during startup recovery"
1000 ),
1001 }
1002 }
1003
1004 let tool_factory =
1005 crate::tools::ToolSurfaceFactory::new(base_tools, tools_with_task, tools);
1006
1007 let session_repo = bamboo_engine::SessionRepository::new(
1008 sessions.clone(),
1009 storage.clone(),
1010 persistence.clone(),
1011 );
1012
1013 let connect_manager = Arc::new(
1018 build_connect_manager(
1019 agent.clone(),
1020 tool_factory.get(crate::tools::ToolSurface::Root),
1021 session_repo.clone(),
1022 agent_runners.clone(),
1023 session_event_senders.clone(),
1024 Some(account_sink.inbox()),
1025 Some(data_dir.clone()),
1026 config.clone(),
1027 provider_registry.clone(),
1028 permission_checker.clone(),
1029 project_store.clone(),
1030 workspace_resolver.clone(),
1031 )
1032 .await
1033 .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?,
1034 );
1035
1036 let child_adapter = Arc::new(crate::tools::ChildSessionAdapter {
1043 session_store: session_store.clone(),
1044 storage: storage.clone(),
1045 persistence: persistence.clone(),
1046 session_messenger: Some(session_messenger.clone()),
1047 scheduler: spawn_scheduler.clone(),
1048 sessions_cache: sessions.clone(),
1049 agent_runners: agent_runners.clone(),
1050 session_event_senders: session_event_senders.clone(),
1051 subagent_model_resolver: None,
1052 config: config.clone(),
1053 project_store: Some(project_store.clone()),
1054 workspace_resolver: workspace_resolver.clone(),
1055 parent_wait_slots: Arc::new(dashmap::DashMap::new()),
1056 });
1057 let guardian_spawner: Arc<dyn bamboo_engine::GuardianSpawner> = child_adapter.clone();
1058 child_completion_coordinator
1061 .set_guardian_spawner(guardian_spawner.clone())
1062 .await;
1063
1064 let bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook> =
1068 child_completion_coordinator.clone();
1069
1070 child_completion_coordinator.spawn_child_wait_watchdog();
1076
1077 let boot_reconcile_services_handle = {
1095 let service_manager = service_manager.clone();
1096 let tool_event_router = tool_event_router.clone();
1097 let app_data_dir = bamboo_home_dir.clone();
1098 tokio::spawn(async move {
1099 crate::plugin_installer::boot_reconcile_services(
1100 &app_data_dir,
1101 &service_manager,
1102 &tool_event_router,
1103 )
1104 .await;
1105 })
1106 };
1107
1108 let (config_watcher, config_live_health, mcp_config_live_health) =
1109 super::config_runtime::ConfigWatcherRuntime::start(
1110 bamboo_home_dir.clone(),
1111 config.clone(),
1112 config_facade.clone(),
1113 config_io_lock.clone(),
1114 provider_registry.clone(),
1115 provider_lock.clone(),
1116 mcp_manager.clone(),
1117 account_sink.clone(),
1118 );
1119 Ok(Self {
1120 app_data_dir: bamboo_home_dir,
1121 memory_store,
1122 tool_event_publisher,
1123 tool_event_router,
1124 config,
1125 config_facade,
1126 config_io_lock,
1127 config_live_health,
1128 mcp_config_live_health,
1129 config_watcher,
1130 project_resource_watcher,
1131 credential_store,
1132 fabric_deployer,
1133 embedded_broker,
1134 health_monitor,
1135 provider: provider_lock,
1136 provider_handle,
1137 sessions,
1138 storage,
1139 session_store,
1140 session_create_operations,
1141 mutation_idempotency,
1142 project_store,
1143 project_context_resolver,
1144 workspace_resolver,
1145 session_repo,
1146 persistence,
1147 session_inbox,
1148 session_activation_router,
1149 session_messenger,
1150 spawn_scheduler,
1151 child_completion_coordinator,
1152 guardian_spawner,
1153 bash_resume_hook,
1154 schedule_store,
1155 schedule_manager,
1156 connect_manager,
1157 tool_factory,
1158 permission_checker,
1159 permission_section,
1160 permission_io_lock,
1161 approval_registry,
1162 notification_service,
1163 session_watchers,
1164 cancel_tokens: Arc::new(RwLock::new(HashMap::new())),
1165 mcp_proxy_shutdown,
1166 skill_manager,
1167 workflow_runs,
1168 mcp_manager,
1169 service_manager,
1170 boot_reconcile_services_handle: tokio::sync::Mutex::new(Some(
1171 boot_reconcile_services_handle,
1172 )),
1173 metrics_service,
1174 agent_runners,
1175 execute_startups: Arc::new(std::sync::Mutex::new(HashMap::new())),
1176 session_event_senders,
1177 account_sink,
1178 process_registry,
1179 metrics_bus: None, agent,
1181 provider_registry,
1182 provider_router,
1183 model_catalog,
1184 title_gen_in_flight: Arc::new(dashmap::DashSet::new()),
1185 pairing_codes: Arc::new(dashmap::DashMap::new()),
1186 pairing_code_guard: Arc::new(crate::handlers::settings::PairingCodeGuard::default()),
1187 root_password_guard: Arc::new(crate::handlers::settings::RootPasswordGuard::default()),
1188 codex_run_tokens,
1189 })
1191 }
1192}
1193
1194pub struct EmbeddedBroker {
1197 task: tokio::task::JoinHandle<()>,
1198 gc_task: tokio::task::JoinHandle<()>,
1199}
1200
1201impl Drop for EmbeddedBroker {
1202 fn drop(&mut self) {
1203 self.task.abort();
1204 self.gc_task.abort();
1205 }
1206}
1207
1208pub struct HealthMonitor(tokio::task::JoinHandle<()>);
1212
1213impl Drop for HealthMonitor {
1214 fn drop(&mut self) {
1215 self.0.abort();
1216 }
1217}
1218
1219async fn maybe_embed_broker(
1231 config: &mut bamboo_llm::Config,
1232 data_dir: &std::path::Path,
1233) -> Option<EmbeddedBroker> {
1234 if let Some(external) = load_external_broker(data_dir) {
1239 let endpoint = external.endpoint.trim().to_string();
1240 if broker_endpoint_reachable(&endpoint).await {
1241 tracing::info!(%endpoint, "using external broker from broker.json");
1242 config.subagents_mut().broker = Some(external);
1243 return None;
1244 }
1245 tracing::warn!(
1246 %endpoint,
1247 "broker.json endpoint is unreachable — embedding a fresh in-process broker instead"
1248 );
1249 }
1250
1251 let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
1252 Ok(l) => l,
1253 Err(e) => {
1254 tracing::warn!("embedded broker: bind failed, sub-agent dispatch disabled: {e}");
1255 return None;
1256 }
1257 };
1258 let port = match listener.local_addr() {
1259 Ok(a) => a.port(),
1260 Err(e) => {
1261 tracing::warn!("embedded broker: local_addr failed: {e}");
1262 return None;
1263 }
1264 };
1265 let token = uuid::Uuid::new_v4().simple().to_string();
1266 let root = data_dir.join("broker");
1267 let core = Arc::new(bamboo_broker::BrokerCore::new(root));
1268 let gc_task = core
1271 .clone()
1272 .spawn_mailbox_gc(std::time::Duration::from_secs(300));
1273 let server = Arc::new(bamboo_broker::BrokerServer::new(core, token.clone()));
1274
1275 let task = tokio::spawn(async move {
1276 if let Err(e) = server.serve(listener).await {
1277 tracing::error!("embedded broker serve loop ended: {e}");
1278 }
1279 });
1280
1281 config.subagents_mut().broker = Some(bamboo_config::BrokerClientConfig {
1284 endpoint: format!("ws://127.0.0.1:{port}"),
1285 token,
1286 token_encrypted: None,
1287 credential_ref: None,
1288 configured: false,
1289 });
1290 tracing::info!(port, "embedded mailbox bus (broker) started in-process");
1291 Some(EmbeddedBroker { task, gc_task })
1292}
1293
1294fn load_external_broker(data_dir: &std::path::Path) -> Option<bamboo_config::BrokerClientConfig> {
1303 let path = data_dir.join("broker.json");
1304 if let Err(error) = bamboo_config::migrate_external_broker_credentials(data_dir)
1305 .and_then(|_| bamboo_config::ensure_provider_mcp_migration_ready(data_dir))
1306 {
1307 tracing::warn!(error = %error, "external broker credential migration unavailable");
1308 return None;
1309 }
1310 let bytes = std::fs::read(&path).ok()?;
1311 match parse_external_broker_snapshot(&bytes, data_dir) {
1312 Ok(cfg) => Some(cfg),
1313 Err(error) => {
1314 tracing::warn!(?path, %error, "broker.json is unavailable");
1315 None
1316 }
1317 }
1318}
1319
1320fn parse_external_broker_snapshot(
1321 bytes: &[u8],
1322 data_dir: &std::path::Path,
1323) -> bamboo_config::ConfigStoreResult<bamboo_config::BrokerClientConfig> {
1324 let mut cfg: bamboo_config::BrokerClientConfig = serde_json::from_slice(bytes)?;
1325 if !cfg.token.trim().is_empty()
1326 || cfg
1327 .token_encrypted
1328 .as_deref()
1329 .is_some_and(|value| !value.trim().is_empty())
1330 {
1331 return Err(bamboo_config::ConfigStoreError::Validation(
1332 "legacy broker credential appeared after migration".to_string(),
1333 ));
1334 }
1335 if cfg.endpoint.trim().is_empty() {
1336 return Err(bamboo_config::ConfigStoreError::Validation(
1337 "broker endpoint is empty".to_string(),
1338 ));
1339 }
1340 cfg.hydrate_credential_from_store(data_dir)?;
1341 Ok(cfg)
1342}
1343
1344async fn broker_endpoint_reachable(endpoint: &str) -> bool {
1349 let host_port = endpoint
1350 .trim()
1351 .trim_start_matches("wss://")
1352 .trim_start_matches("ws://")
1353 .split('/')
1354 .next()
1355 .unwrap_or("");
1356 if host_port.is_empty() {
1357 return false;
1358 }
1359 matches!(
1360 tokio::time::timeout(
1361 std::time::Duration::from_millis(500),
1362 tokio::net::TcpStream::connect(host_port),
1363 )
1364 .await,
1365 Ok(Ok(_))
1366 )
1367}
1368
1369#[cfg(test)]
1370mod fabric_boot_reconcile_tests {
1371 use super::*;
1372 use bamboo_config::cluster_fabric::{
1373 DeployProfile, Node, NodePlacement, NodeState, NodeStatus, TrustLevel,
1374 };
1375 use bamboo_config::{ClusterNodeCredentialIntents, ConfigFacade};
1376 use std::collections::BTreeMap;
1377
1378 #[tokio::test]
1379 async fn restart_reconcile_keeps_runtime_process_facade_and_disk_on_one_revision() {
1380 let _key = bamboo_config::encryption::set_test_encryption_key([0x7b; 32]);
1381 let dir = tempfile::tempdir().unwrap();
1382 let first = AppState::new(dir.path().to_path_buf()).await.unwrap();
1383 first
1384 .update_cluster_fabric_credentials(
1385 0,
1386 BTreeMap::from([(
1387 "boot-node".to_string(),
1388 ClusterNodeCredentialIntents::clear_all(),
1389 )]),
1390 |config| {
1391 config.cluster_fabric.nodes.push(Node {
1392 id: "boot-node".to_string(),
1393 label: "boot-node".to_string(),
1394 placement: NodePlacement::Local,
1395 trust_level: TrustLevel::Trusted,
1396 deploy: DeployProfile::default(),
1397 state: Some(NodeState {
1398 status: NodeStatus::Running,
1399 worker_id: Some("stale-worker".to_string()),
1400 ..Default::default()
1401 }),
1402 enabled: true,
1403 });
1404 Ok(())
1405 },
1406 )
1407 .await
1408 .unwrap();
1409 assert_eq!(
1410 first
1411 .config_facade
1412 .as_ref()
1413 .unwrap()
1414 .registry()
1415 .cluster_fabric
1416 .snapshot()
1417 .revision,
1418 1
1419 );
1420 drop(first);
1421
1422 let restarted = AppState::new(dir.path().to_path_buf()).await.unwrap();
1423 let process_snapshot = restarted
1424 .config_facade
1425 .as_ref()
1426 .unwrap()
1427 .registry()
1428 .cluster_fabric
1429 .snapshot();
1430 assert_eq!(process_snapshot.revision, 2);
1431 assert_eq!(
1432 process_snapshot
1433 .data
1434 .0
1435 .node("boot-node")
1436 .unwrap()
1437 .state
1438 .as_ref()
1439 .unwrap()
1440 .status,
1441 NodeStatus::Unreachable
1442 );
1443 assert_eq!(
1444 restarted
1445 .config
1446 .read()
1447 .await
1448 .cluster_fabric
1449 .node("boot-node")
1450 .unwrap()
1451 .state
1452 .as_ref()
1453 .unwrap()
1454 .status,
1455 NodeStatus::Unreachable
1456 );
1457 let reopened = ConfigFacade::open(dir.path()).unwrap();
1458 assert_eq!(reopened.registry().cluster_fabric.snapshot().revision, 2);
1459 assert_eq!(
1460 reopened
1461 .effective_config()
1462 .cluster_fabric
1463 .node("boot-node")
1464 .unwrap()
1465 .state
1466 .as_ref()
1467 .unwrap()
1468 .status,
1469 NodeStatus::Unreachable
1470 );
1471 }
1472}
1473
1474#[cfg(test)]
1475mod broker_embed_tests {
1476 use super::{broker_endpoint_reachable, load_external_broker, parse_external_broker_snapshot};
1477
1478 #[test]
1479 fn load_external_broker_reads_broker_json_not_config() {
1480 let _key = bamboo_config::encryption::set_test_encryption_key([0x57; 32]);
1481 let dir = tempfile::tempdir().unwrap();
1482 assert!(load_external_broker(dir.path()).is_none());
1484
1485 std::fs::write(
1487 dir.path().join("broker.json"),
1488 r#"{ "endpoint": "wss://broker.example:9600", "token": "t" }"#,
1489 )
1490 .unwrap();
1491 let got = load_external_broker(dir.path()).expect("parsed");
1492 assert_eq!(got.endpoint, "wss://broker.example:9600");
1493 assert_eq!(got.token, "t");
1494 assert_eq!(
1495 got.credential_ref.as_ref().unwrap().as_str(),
1496 "broker.external.bearer_token"
1497 );
1498 assert!(got.configured);
1499 let durable = std::fs::read_to_string(dir.path().join("broker.json")).unwrap();
1500 assert!(!durable.contains("\"token\""));
1501 assert!(!durable.contains("token_encrypted"));
1502 assert!(!durable.contains("\"t\""));
1503
1504 std::fs::write(dir.path().join("broker.json"), r#"{ "endpoint": " " }"#).unwrap();
1506 assert!(load_external_broker(dir.path()).is_none());
1507
1508 std::fs::write(dir.path().join("broker.json"), "not json").unwrap();
1510 assert!(load_external_broker(dir.path()).is_none());
1511 }
1512
1513 #[test]
1514 fn external_broker_reference_fails_closed_and_tracks_generic_cas_updates() {
1515 let _key = bamboo_config::encryption::set_test_encryption_key([0x58; 32]);
1516 let dir = tempfile::tempdir().unwrap();
1517 let reference =
1518 bamboo_config::credential_ref("broker", "external", "bearer_token").unwrap();
1519 std::fs::write(
1520 dir.path().join("broker.json"),
1521 serde_json::to_vec_pretty(&serde_json::json!({
1522 "endpoint": "wss://broker.example:9600",
1523 "credential_ref": reference,
1524 "configured": true,
1525 }))
1526 .unwrap(),
1527 )
1528 .unwrap();
1529
1530 assert!(load_external_broker(dir.path()).is_none());
1531 std::fs::write(
1532 dir.path().join("broker.json"),
1533 serde_json::to_vec_pretty(&serde_json::json!({
1534 "endpoint": "wss://broker.example:9600",
1535 "credential_ref": reference,
1536 "configured": false,
1537 }))
1538 .unwrap(),
1539 )
1540 .unwrap();
1541 assert!(load_external_broker(dir.path()).is_none());
1542 let store = bamboo_config::CredentialStore::open(dir.path());
1543 store
1544 .replace(
1545 reference.clone(),
1546 "replacement-token",
1547 bamboo_config::CredentialSource::User,
1548 0,
1549 )
1550 .unwrap();
1551 let hydrated = load_external_broker(dir.path()).unwrap();
1552 assert_eq!(hydrated.token, "replacement-token");
1553 store.clear(&reference, 1).unwrap();
1554 assert!(load_external_broker(dir.path()).is_none());
1555 }
1556
1557 #[test]
1558 fn broker_snapshot_with_late_legacy_secret_fails_closed() {
1559 let _key = bamboo_config::encryption::set_test_encryption_key([0x59; 32]);
1560 let dir = tempfile::tempdir().unwrap();
1561 let error = parse_external_broker_snapshot(
1562 br#"{
1563 "endpoint": "wss://broker.example:9600",
1564 "token": "late-legacy-token"
1565 }"#,
1566 dir.path(),
1567 )
1568 .unwrap_err();
1569 assert!(error.to_string().contains("appeared after migration"));
1570 }
1571
1572 #[tokio::test]
1573 async fn reachability_probe_distinguishes_live_from_dead() {
1574 let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1577 let dead = l.local_addr().unwrap();
1578 drop(l);
1579 assert!(!broker_endpoint_reachable(&format!("ws://{dead}")).await);
1580
1581 assert!(!broker_endpoint_reachable("").await);
1583 assert!(!broker_endpoint_reachable("ws://").await);
1584
1585 let live = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1587 let addr = live.local_addr().unwrap();
1588 assert!(broker_endpoint_reachable(&format!("ws://{addr}/stream")).await);
1589 }
1590}