Skip to main content

bamboo_server/app_state/
builder.rs

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    /// Create unified app state with direct provider access
48    ///
49    /// This eliminates the proxy pattern where we created an AgentAppState
50    /// that called back to web_service via HTTP. Now we have direct provider access.
51    ///
52    /// # Arguments
53    ///
54    /// * `bamboo_home_dir` - Bamboo home directory containing all application data.
55    ///   This is the root directory (e.g., `${HOME}/.bamboo`) that contains:
56    ///   - config.json: Configuration file
57    ///   - sessions/: Conversation history
58    ///   - skills/: Skill definitions
59    ///   - workflows/: Workflow definitions
60    ///   - cache/: Cached data
61    ///   - runtime/: Runtime files
62    ///   - workspaces/: Default per-session workspace dirs (issue #217) — a
63    ///     session with no configured/explicit workspace gets
64    ///     `workspaces/{session_id}` here instead of the server process's
65    ///     cwd. Overridable via `BAMBOO_WORKSPACE_ROOT`.
66    ///   - subagents/: Local actor sub-agent fabric discovery + isolated
67    ///     per-child storage (issue #217) — replaces the old
68    ///     `env::temp_dir()/bamboo-subagents` default.
69    ///
70    /// # Returns
71    ///
72    /// A fully initialized AppState with all components ready for use.
73    /// # Example
74    ///
75    /// ```rust,no_run
76    /// use bamboo_server::app_state::AppState;
77    /// use std::path::PathBuf;
78    ///
79    /// #[tokio::main]
80    /// async fn main() {
81    ///     let state = AppState::new(PathBuf::from("/path/to/bamboo-data-dir"))
82    ///         .await
83    ///         .expect("failed to initialize app state");
84    ///     let provider = state.get_provider().await;
85    ///     let _models = provider.list_models().await.ok();
86    /// }
87    /// ```
88    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    /// Create a server state with one explicit Jiandu store shared by the tool,
94    /// prompt, Dream, gardener, metrics, and maintenance paths.
95    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        // Ensure all helpers that rely on `core::paths::bamboo_dir()` see the same
100        // directory as the server runtime.
101        bamboo_config::paths::init_bamboo_dir(bamboo_home_dir.clone());
102
103        // Complete the recoverable legacy split before any runtime component
104        // reads configuration, then keep this facade as the process authority.
105        // A malformed legacy document or an unrecoverable pending legacy
106        // transaction must not make the server exit: retain the compatibility
107        // loader's recovered/LKG view for this process so the recovery API can
108        // still repair or confirm it. Healthy layouts never take that fallback.
109        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        // Loud, unmissable startup signal: `plugin_trust.enforcement: off`
139        // silently affects EVERY future `bamboo plugin install`/`update`
140        // (URL sources skip the host allowlist, signature, and checksum
141        // layers with no per-install flag needed — see
142        // `bamboo_server::plugin_source`'s module docs), not just one
143        // command invocation, so it gets its own warning here at boot in
144        // addition to the per-install warning `fetch_manifest_bundle` logs
145        // for each individual insecure install. The live config-apply paths
146        // (`update_config`/`replace_config`) emit the SAME warning on a flip
147        // to `Off`, so no trigger — boot, `bamboo config set`, or an HTTP
148        // config PATCH — can relax it silently.
149        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    /// Create unified app state with a specific provider
198    ///
199    /// Allows injecting a custom LLM provider instead of creating
200    /// one from configuration. Useful for testing and custom deployments.
201    ///
202    /// # Arguments
203    ///
204    /// * `bamboo_home_dir` - Bamboo home directory containing all application data
205    /// * `config` - Application configuration
206    /// * `provider` - Pre-configured LLM provider implementation
207    ///
208    /// # Returns
209    ///
210    /// A fully initialized AppState with the provided provider.
211    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    /// Create a server state with an instance-local tool-event publisher.
229    /// Existing constructors retain the no-op publisher and unchanged behavior.
230    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        // Wire the configured-default-workspace resolver into agent-core. This keeps
257        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        // In-memory session cache (shared across handlers and background jobs).
301        let sessions: bamboo_engine::SessionCache = Arc::default();
302
303        // Embed the mailbox bus (broker) in-process unless an external one is
304        // configured. Mutates `config.subagents.broker` to point at the loopback
305        // bus BEFORE it is wrapped/read downstream, so ask_agent / deploy_agent /
306        // cluster all wire to it — and a standalone `bamboo broker serve` is no
307        // longer required for sub-agent dispatch. (Foundation for routing local
308        // actors onto the bus.)
309        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        // Build one coherent configured-default/root provider pair. The process
315        // globals below intentionally remain first-registration-wins, while
316        // this AppState retains the same pair for instance-scoped Project
317        // preview and post-persistence publication (#717).
318        //
319        // The default closure reads the LIVE in-memory config, not a fresh
320        // disk-reading Config::new() (#38). `try_read` never blocks (the
321        // resolver is called from sync code); on write-lock contention it
322        // returns the last successful value rather than transiently falling
323        // back to another source.
324        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        // Issue #217: the root provider re-reads operator env policy on every
341        // call. Its no-override fallback is this AppState's own data directory,
342        // not the unrelated process-global bamboo_dir that another test state
343        // may have registered first.
344        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        // Service input and ToolEvent routing share one AppState-owned
389        // lifecycle. The router is built before the tool surfaces so every
390        // executor receives the production publisher from its first call;
391        // it remains inert until install/boot reconciliation declares sinks.
392        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        // NOTE: the idle-eviction sweep for `agent_runners` is spawned below,
424        // once `session_event_senders` also exists, so it can drop both maps'
425        // entries for a completed session together (issue #346).
426
427        let process_registry = Arc::new(ProcessRegistry::new());
428        let (provider_lock, provider_handle) = build_provider_handles(provider);
429
430        // Initialize multi-provider registry (for features.provider_model_ref).
431        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        // Long-lived session event senders map (UI subscriptions + background tasks).
461        // Declared before `build_base_tools` (moved up from its original spot below)
462        // because the `notify` tool overlaid there needs it to broadcast onto a
463        // session's live channel — see `app_state::tools::build_base_tools`.
464        let session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>> =
465            Arc::new(RwLock::new(HashMap::new()));
466
467        // Shared bundle of always-on notification relay deps (see
468        // `session_events::NotificationRelayDeps`). Built once and cloned into
469        // every entry point that starts a relay directly at execution time —
470        // the schedule manager, the root child-session adapter, and the
471        // guardian child-session adapter below — so they can never drift.
472        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        // The `ledger` tool needs the schedule store (built further down) to
480        // sync reminders; hand it a late-bound bridge now and bind it below.
481        let ledger_schedule_bridge =
482            Arc::new(crate::schedule_app::LateBoundLedgerBridge::default());
483
484        // The runtime and skill tools must share one cache-aware coordinator.
485        // Session setup publishes this run's resolved skill allowlist through it
486        // before the model can call load_skill.
487        let session_repo = bamboo_engine::SessionRepository::new(
488            sessions.clone(),
489            storage.clone(),
490            persistence.clone(),
491        );
492
493        // Account-scoped durable change feed. It is initialized before the
494        // Project tool surface so non-HTTP Project mutations publish the same
495        // replayable events as the REST API.
496        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        // The workflow engine executes against the base tool surface. The
534        // caller-facing workflow_run tool is overlaid onto the root surface
535        // later, preventing a workflow from recursively dispatching itself.
536        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        // Idle-evict completed runners together with their paired session event
547        // senders (issue #346). Spawned here (not next to `agent_runners`) so it
548        // owns handles to both maps.
549        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        // Bridge both instruction and orchestration catalog transitions onto
558        // the same durable account feed used by SSE and v2 WebSocket clients.
559        // The public Workflow Library unifies both kinds, so either namespace
560        // must invalidate clients. Catalog events are account-scoped (no
561        // session id).
562        {
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        // Sub-agents are full agents with the full toolset (no per-role tool
616        // trimming): the child tool surface is the plain base tools.
617        let child_tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = base_tools.clone();
618
619        // Unified agent runtime (shared resources for all execution paths).
620        // default_tools = base_tools (builtin + MCP + memory + skills) as a safe fallback.
621        // Interactive execution paths pass an explicit tool surface override:
622        // root sessions use ToolSurface::Root; child sessions use ToolSurface::Child.
623        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        // Initialize sub-session spawn scheduler (async background jobs).
673        let config_snapshot = config.read().await.clone();
674
675        // When a broker is configured, run the MCP proxy service under a
676        // supervisor (issue #47): deployed workers forward their (host-bound)
677        // MCP tool calls here, and we execute them against this orchestrator's
678        // real MCP servers (single MCP host). The supervisor restarts the proxy
679        // with bounded backoff after a transient WebSocket drop instead of
680        // permanently disabling proxied tools for the worker's lifetime.
681        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                // Build the orchestrator-side per-role MCP tool allowlist from
692                // config (issue #54). Enforcement lives HERE — never in the
693                // worker-facing `McpProxyConfig` a deployed worker receives —
694                // because a worker self-declaring its own allowlist would be
695                // insecure (it could simply claim to be unrestricted). Validate
696                // configured tool names against THIS backend's real, live tool
697                // set so a typo in `config.json` is surfaced at boot instead of
698                // silently granting nothing for the intended tool.
699                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                    // Default behavior unchanged (issue #54 item 5): no policy
707                    // configured -> every role sees/can call every proxied
708                    // tool. Logged once here at boot (not per-request) so an
709                    // operator can tell the restriction is opt-in rather than
710                    // silently absent.
711                    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                // NOTE: `init_mcp_manager` connects MCP servers in a background
718                // task so the HTTP API stays responsive at boot — this backend's
719                // `list_tools()` can legitimately be empty here if that task
720                // hasn't finished yet. `from_config` treats an empty set as "skip
721                // tool-name validation" rather than flagging every configured
722                // tool as an unknown typo, so warn separately here when that
723                // race means validation was effectively skipped.
724                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        // Bind the ledger's reminder bridge now that the schedule store exists.
806        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        // Background memory "gardener": opt-in blob remediation + near-duplicate
843        // consolidation. No-op cost unless `memory.gardener_enabled` /
844        // `memory.dedup_gardener_enabled` is set; an empty prefilter makes zero LLM calls.
845        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        // Background ledger gardener: expiry + record↔schedule reconciliation
858        // are deterministic and free; distillation uses the background model
859        // and no-ops without one. The bridge handle is already bound above.
860        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                            &registry,
890                            &subagent_type,
891                        )
892                    })
893                },
894            ))
895        };
896
897        // Config-write io-lock + the shared Remote Cluster Fabric deploy engine.
898        // The engine is built once and shared by the HTTP handlers (via AppState)
899        // and the `cluster` agent tool, so both use ONE worker registry.
900        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        // Cluster health monitor: periodically probe deployed workers on the bus and
935        // flip node status live (Running↔Unreachable) + auto-recover. Server-scoped
936        // — it runs under BOTH the embedded and an external broker (it reads the
937        // broker endpoint lazily each tick), and is aborted when the server drops.
938        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        // Process restart recovery: only backlog covered by its producer's
975        // durable activation watermark requests a run. A child/Bash coordinator
976        // may intentionally stage sibling outcomes while a specific wait remains
977        // armed; admission by itself is not permission to execute.
978        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        // bamboo-connect (#452 / epic #447): drives bamboo sessions from IM
1014        // platforms (Telegram first). Fully inert when `config.connect.platforms`
1015        // is empty — mirrors the schedule manager / notification relay's
1016        // always-constructed-but-only-active-if-configured lifecycle.
1017        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        // Dedicated child-session adapter backing the guardian review spawner.
1037        // The guardian path passes an explicit model (no subagent_type routing)
1038        // and registers its parent wait at the terminal gate (not via the
1039        // adapter's coalescing slots), so a lightweight adapter with no resolver
1040        // and a fresh wait-slot map suffices. `Arc<ChildSessionAdapter>` doubles
1041        // as `Arc<dyn GuardianSpawner>`.
1042        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        // Wire the spawner into the completion coordinator too, so a resumed run
1059        // can re-spawn a guardian to re-review a fix after a reject verdict.
1060        child_completion_coordinator
1061            .set_guardian_spawner(guardian_spawner.clone())
1062            .await;
1063
1064        // The completion coordinator doubles as the bash self-resume hook
1065        // (issue #84 Phase 2b): it polls the live shell registry and resumes a
1066        // session once all its background bash shells finish.
1067        let bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook> =
1068            child_completion_coordinator.clone();
1069
1070        // Child-wait watchdog (issue #546): boot-time reconciliation of
1071        // children orphaned by a restart, then a periodic heartbeat sweep that
1072        // backstops every lost child→parent wake (panicked child task, dead
1073        // spawn scheduler, clobbered resume, expired wait lease, ...). Spawned
1074        // AFTER `set_root_tools` above so a boot-time parent resume can spawn.
1075        child_completion_coordinator.spawn_child_wait_watchdog();
1076
1077        // `ServiceManager` (issue #479 / epic #477 prereq): supervises
1078        // long-running "service" plugins. Always constructed — fully inert
1079        // until a plugin install or the boot-time reconcile below starts
1080        // something — mirrors `mcp_manager`/`connect_manager`'s
1081        // always-alive lifecycle.
1082        // Backgrounded (mirrors `init_mcp_manager`'s background MCP
1083        // bootstrap): a service that `installed.json` says should be
1084        // running but isn't (the previous `bamboo serve` process, if
1085        // any, died with everything it supervised) is started fresh.
1086        //
1087        // The `JoinHandle` is kept (not discarded) purely so tests can
1088        // deterministically wait it out via
1089        // `AppState::wait_for_boot_reconcile_services`. The pass shares the
1090        // plugin operation lock with install/update/uninstall, so its
1091        // manifest/provenance generation cannot race a newer mutation; see
1092        // that method's doc comment and issue #486. Production code never
1093        // awaits it; server startup is never blocked on plugin service spawns.
1094        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, // Will be set by server if needed
1180            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            // remote-actor P2a (#181): empty in-memory agent registry.
1190        })
1191    }
1192}
1193
1194/// A handle to the in-process mailbox bus (broker) so it can be shut down with
1195/// the server. Dropping it aborts the serve task.
1196pub 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
1208/// Server-scoped handle to the cluster health monitor. Kept separate from
1209/// [`EmbeddedBroker`] so the monitor runs under BOTH the embedded broker and an
1210/// external (`broker.json`) one; aborts the sweep on drop.
1211pub struct HealthMonitor(tokio::task::JoinHandle<()>);
1212
1213impl Drop for HealthMonitor {
1214    fn drop(&mut self) {
1215        self.0.abort();
1216    }
1217}
1218
1219/// Start an in-process broker on `127.0.0.1:<auto>` and point
1220/// `config.subagents.broker` (RUNTIME-ONLY, `#[serde(skip)]`) at it — UNLESS a
1221/// user-managed external broker is configured in `<data_dir>/broker.json` and
1222/// reachable (then use that). Returns `None` when an external broker is used or
1223/// the bind fails (sub-agent dispatch then degrades exactly as before).
1224///
1225/// The broker's endpoint deliberately lives in EITHER the in-memory config (for
1226/// the embedded case, regenerated each boot) or its own `broker.json` (for the
1227/// external case) — NEVER in `config.json`. That is what stops a prior run's
1228/// ephemeral auto-port from leaking into the user's config and being dialed dead
1229/// on the next boot (every sub-agent + the MCP proxy would hit "connect refused").
1230async fn maybe_embed_broker(
1231    config: &mut bamboo_llm::Config,
1232    data_dir: &std::path::Path,
1233) -> Option<EmbeddedBroker> {
1234    // A user-managed EXTERNAL broker (multi-host / shared standalone bus) lives in
1235    // its OWN file, `<data_dir>/broker.json` — separate from config.json. Honour
1236    // it ONLY if actually REACHABLE; a dead endpoint (standalone broker not up)
1237    // falls through to a fresh in-process broker so dispatch still works.
1238    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    // Reclaim orphan mailbox dirs (one-shot parent links, killed pool workers)
1269    // every 5 min so `<data>/broker/mailboxes/` doesn't grow unbounded.
1270    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    // Set the endpoint IN MEMORY ONLY — `subagents.broker` is `#[serde(skip)]`, so
1282    // this ephemeral loopback port is regenerated every boot and never touches disk.
1283    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
1294/// Load a user-managed EXTERNAL broker from `<data_dir>/broker.json`, if present.
1295/// This file is the SEPARATE, persisted home for a standalone/remote broker —
1296/// deliberately NOT `config.json`, so the embedded broker's ephemeral runtime
1297/// port can never leak into the user's config (the stale-dead-port bug). An
1298/// absent file or a parse error yields `None` (embed a fresh in-process broker).
1299///
1300/// Durable format carries only endpoint plus credential metadata. Legacy
1301/// plaintext/ciphertext tokens are migrated before this reader proceeds.
1302fn 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
1344/// Best-effort TCP reachability probe of a `ws[s]://host:port[/path]` broker
1345/// endpoint. Used to tell a LIVE external/standalone broker (keep) apart from a
1346/// DEAD persisted endpoint (a prior run's embedded auto-port that leaked into
1347/// config.json — re-embed). A short timeout keeps boot fast when it's dead.
1348async 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        // Absent file ⇒ None (embed a fresh in-process broker).
1483        assert!(load_external_broker(dir.path()).is_none());
1484
1485        // A well-formed broker.json ⇒ parsed external broker.
1486        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        // Empty endpoint ⇒ ignored (treated as absent).
1505        std::fs::write(dir.path().join("broker.json"), r#"{ "endpoint": "  " }"#).unwrap();
1506        assert!(load_external_broker(dir.path()).is_none());
1507
1508        // Garbage ⇒ ignored, never panics.
1509        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        // Bound-then-dropped port: nothing listening ⇒ unreachable (a stale
1575        // persisted embedded endpoint must NOT be trusted).
1576        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        // Empty / malformed ⇒ never trusted.
1582        assert!(!broker_endpoint_reachable("").await);
1583        assert!(!broker_endpoint_reachable("ws://").await);
1584
1585        // A LIVE listener (with a path) ⇒ reachable (a real external broker is kept).
1586        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}