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::tools::OptionalSubagentModelResolver;
9use bamboo_agent_core::storage::Storage;
10
11impl AppState {
12    /// Create unified app state with direct provider access
13    ///
14    /// This eliminates the proxy pattern where we created an AgentAppState
15    /// that called back to web_service via HTTP. Now we have direct provider access.
16    ///
17    /// # Arguments
18    ///
19    /// * `bamboo_home_dir` - Bamboo home directory containing all application data.
20    ///   This is the root directory (e.g., `${HOME}/.bamboo`) that contains:
21    ///   - config.json: Configuration file
22    ///   - sessions/: Conversation history
23    ///   - skills/: Skill definitions
24    ///   - workflows/: Workflow definitions
25    ///   - cache/: Cached data
26    ///   - runtime/: Runtime files
27    ///   - workspaces/: Default per-session workspace dirs (issue #217) — a
28    ///     session with no configured/explicit workspace gets
29    ///     `workspaces/{session_id}` here instead of the server process's
30    ///     cwd. Overridable via `BAMBOO_WORKSPACE_ROOT`.
31    ///   - subagents/: Local actor sub-agent fabric discovery + isolated
32    ///     per-child storage (issue #217) — replaces the old
33    ///     `env::temp_dir()/bamboo-subagents` default.
34    ///
35    /// # Returns
36    ///
37    /// A fully initialized AppState with all components ready for use.
38    /// # Example
39    ///
40    /// ```rust,no_run
41    /// use bamboo_server::app_state::AppState;
42    /// use std::path::PathBuf;
43    ///
44    /// #[tokio::main]
45    /// async fn main() {
46    ///     let state = AppState::new(PathBuf::from("/path/to/bamboo-data-dir"))
47    ///         .await
48    ///         .expect("failed to initialize app state");
49    ///     let provider = state.get_provider().await;
50    ///     let _models = provider.list_models().await.ok();
51    /// }
52    /// ```
53    pub async fn new(bamboo_home_dir: PathBuf) -> Result<Self, AppError> {
54        // Ensure all helpers that rely on `core::paths::bamboo_dir()` see the same
55        // directory as the server runtime.
56        bamboo_config::paths::init_bamboo_dir(bamboo_home_dir.clone());
57
58        // Complete the recoverable legacy split before any runtime component
59        // reads configuration, then keep this facade as the process authority.
60        // A malformed legacy document or an unrecoverable pending legacy
61        // transaction must not make the server exit: retain the compatibility
62        // loader's recovered/LKG view for this process so the recovery API can
63        // still repair or confirm it. Healthy layouts never take that fallback.
64        let (config, config_facade) = match bamboo_config::ConfigFacade::open_or_migrate(
65            &bamboo_home_dir,
66        ) {
67            Ok(facade) => {
68                let facade = Arc::new(facade);
69                let config =
70                    super::config_runtime::load_facade_effective_config(&facade, &bamboo_home_dir);
71                (config, Some(facade))
72            }
73            Err(error) => {
74                if bamboo_config::modular_authority_boundary_present(&bamboo_home_dir)
75                    .unwrap_or(true)
76                {
77                    return Err(AppError::InternalError(anyhow::anyhow!(
78                        "modular configuration authority is unavailable: {error}"
79                    )));
80                }
81                tracing::warn!(
82                    error = %error,
83                    "modular configuration facade is unavailable; retaining recovered legacy authority"
84                );
85                (
86                    Config::from_data_dir_without_publish(Some(bamboo_home_dir.clone())),
87                    None,
88                )
89            }
90        };
91        config.publish_env_vars();
92
93        // Loud, unmissable startup signal: `plugin_trust.enforcement: off`
94        // silently affects EVERY future `bamboo plugin install`/`update`
95        // (URL sources skip the host allowlist, signature, and checksum
96        // layers with no per-install flag needed — see
97        // `bamboo_server::plugin_source`'s module docs), not just one
98        // command invocation, so it gets its own warning here at boot in
99        // addition to the per-install warning `fetch_manifest_bundle` logs
100        // for each individual insecure install. The live config-apply paths
101        // (`update_config`/`replace_config`) emit the SAME warning on a flip
102        // to `Off`, so no trigger — boot, `bamboo config set`, or an HTTP
103        // config PATCH — can relax it silently.
104        if config.plugin_trust.enforcement_is_off() {
105            super::config_runtime::warn_plugin_trust_enforcement_off();
106        }
107
108        let provider_registry =
109            match bamboo_llm::ProviderRegistry::from_config(&config, bamboo_home_dir.clone()).await
110            {
111                Ok(registry) => Arc::new(registry),
112                Err(e) => {
113                    tracing::error!("Failed to create provider registry: {}", e);
114                    Arc::new(
115                        bamboo_llm::ProviderRegistry::from_config(
116                            &Config::default(),
117                            bamboo_home_dir.clone(),
118                        )
119                        .await
120                        .expect("Cannot create even an empty provider registry"),
121                    )
122                }
123            };
124
125        let provider = provider_registry.get_default().unwrap_or_else(|| {
126            let default_provider_name = provider_registry.default_provider_name();
127            let message = if config.has_provider_instances() {
128                format!(
129                    "Default provider instance '{}' is not available or failed to initialize",
130                    default_provider_name
131                )
132            } else {
133                format!(
134                    "Provider '{}' is not available or failed to initialize",
135                    config.provider
136                )
137            };
138            Arc::new(UnconfiguredProvider { message }) as Arc<dyn LLMProvider>
139        });
140
141        Self::new_with_provider_and_facade(bamboo_home_dir, config, provider, config_facade).await
142    }
143
144    /// Create unified app state with a specific provider
145    ///
146    /// Allows injecting a custom LLM provider instead of creating
147    /// one from configuration. Useful for testing and custom deployments.
148    ///
149    /// # Arguments
150    ///
151    /// * `bamboo_home_dir` - Bamboo home directory containing all application data
152    /// * `config` - Application configuration
153    /// * `provider` - Pre-configured LLM provider implementation
154    ///
155    /// # Returns
156    ///
157    /// A fully initialized AppState with the provided provider.
158    pub async fn new_with_provider(
159        bamboo_home_dir: PathBuf,
160        config: Config,
161        provider: Arc<dyn LLMProvider>,
162    ) -> Result<Self, AppError> {
163        Self::new_with_provider_and_facade(bamboo_home_dir, config, provider, None).await
164    }
165
166    async fn new_with_provider_and_facade(
167        bamboo_home_dir: PathBuf,
168        config: Config,
169        provider: Arc<dyn LLMProvider>,
170        config_facade: Option<Arc<bamboo_config::ConfigFacade>>,
171    ) -> Result<Self, AppError> {
172        // Wire the configured-default-workspace resolver into agent-core. This keeps
173        let data_dir = bamboo_home_dir.clone();
174        let (session_store, storage) = init_storage(&data_dir).await?;
175        let session_create_operations =
176            Arc::new(super::session_create_operations::SessionCreateOperationStore::new(&data_dir));
177        let mutation_idempotency =
178            Arc::new(super::mutation_idempotency::MutationIdempotencyStore::default());
179        match session_create_operations.prune_expired().await {
180            Ok(0) => {}
181            Ok(deleted) => tracing::info!(
182                target: "bamboo.session_create",
183                phase = "retention_cleanup",
184                outcome = "expired_pruned",
185                deleted,
186                "pruned expired session-create operation receipts"
187            ),
188            Err(error) => tracing::warn!(
189                target: "bamboo.session_create",
190                phase = "retention_cleanup",
191                outcome = "cleanup_failed",
192                error = %error,
193                "failed to prune expired session-create operation receipts"
194            ),
195        }
196        let project_store = Arc::new(bamboo_projects::ProjectStore::open(&data_dir).map_err(
197            |error| {
198                AppError::InternalError(anyhow::anyhow!(
199                    "failed to initialize Project registry: {error}"
200                ))
201            },
202        )?);
203        let persistence = Arc::new(LockedSessionStore::new(storage.clone()));
204        let session_inbox: Arc<dyn bamboo_domain::SessionInboxPort> =
205            Arc::new(bamboo_storage::FileSessionInbox::new(
206                session_store.clone(),
207                bamboo_domain::SessionInboxLimits::default(),
208            ));
209        let session_activation_router = bamboo_engine::SessionActivationRouter::new();
210        let session_messenger = Arc::new(bamboo_engine::SessionMessenger::new(
211            storage.clone(),
212            session_inbox.clone(),
213            session_activation_router.clone(),
214        ));
215
216        // In-memory session cache (shared across handlers and background jobs).
217        let sessions: bamboo_engine::SessionCache = Arc::new(dashmap::DashMap::new());
218
219        // Embed the mailbox bus (broker) in-process unless an external one is
220        // configured. Mutates `config.subagents.broker` to point at the loopback
221        // bus BEFORE it is wrapped/read downstream, so ask_agent / deploy_agent /
222        // cluster all wire to it — and a standalone `bamboo broker serve` is no
223        // longer required for sub-agent dispatch. (Foundation for routing local
224        // actors onto the bus.)
225        let mut config = config;
226        let embedded_broker = maybe_embed_broker(&mut config, &data_dir).await;
227
228        let config = Arc::new(RwLock::new(config));
229
230        // Build one coherent configured-default/root provider pair. The process
231        // globals below intentionally remain first-registration-wins, while
232        // this AppState retains the same pair for instance-scoped Project
233        // preview and post-persistence publication (#717).
234        //
235        // The default closure reads the LIVE in-memory config, not a fresh
236        // disk-reading Config::new() (#38). `try_read` never blocks (the
237        // resolver is called from sync code); on write-lock contention it
238        // returns the last successful value rather than transiently falling
239        // back to another source.
240        let live_default_workspace: Arc<dyn Fn() -> Option<PathBuf> + Send + Sync> = {
241            let config_for_workspace = config.clone();
242            let last_known: Arc<std::sync::Mutex<Option<PathBuf>>> =
243                Arc::new(std::sync::Mutex::new(None));
244            Arc::new(move || match config_for_workspace.try_read() {
245                Ok(cfg) => {
246                    let path = cfg.get_default_work_area_path();
247                    if let Ok(mut cache) = last_known.lock() {
248                        *cache = path.clone();
249                    }
250                    path
251                }
252                Err(_) => last_known.lock().ok().and_then(|cache| cache.clone()),
253            })
254        };
255
256        // Issue #217: the root provider re-reads operator env policy on every
257        // call. Its no-override fallback is this AppState's own data directory,
258        // not the unrelated process-global bamboo_dir that another test state
259        // may have registered first.
260        let live_workspace_root: Arc<
261            dyn Fn() -> bamboo_agent_core::workspace_state::WorkspaceRootConfig + Send + Sync,
262        > = {
263            let app_data_dir = data_dir.clone();
264            Arc::new(
265                move || bamboo_agent_core::workspace_state::WorkspaceRootConfig {
266                    root: bamboo_config::paths::resolve_workspace_root_in(&app_data_dir),
267                    confine: bamboo_config::paths::workspace_confinement_enforced(),
268                },
269            )
270        };
271
272        let workspace_resolver = bamboo_agent_core::workspace_state::WorkspaceResolver::new(
273            {
274                let provider = live_default_workspace.clone();
275                move || provider()
276            },
277            {
278                let provider = live_workspace_root.clone();
279                move || provider()
280            },
281        );
282
283        bamboo_agent_core::workspace_state::set_default_workspace_provider(Box::new({
284            let provider = live_default_workspace;
285            move || provider()
286        }));
287        bamboo_agent_core::workspace_state::set_workspace_root_provider(Box::new({
288            let provider = live_workspace_root;
289            move || provider()
290        }));
291
292        let (permission_checker, permission_section) =
293            load_permission_checker(&bamboo_home_dir).await?;
294        let permission_io_lock = Arc::new(tokio::sync::Mutex::new(()));
295        let notification_service = Arc::new(bamboo_notification::NotificationService::new(
296            bamboo_home_dir.join("notification_preferences.json"),
297        ));
298        let session_watchers = super::watchers::SessionWatchers::new();
299        let (mcp_manager, _legacy_mcp_bootstrap) =
300            init_mcp_manager(config.clone(), &bamboo_home_dir);
301        let skill_manager = init_skill_manager(&data_dir).await;
302        let metrics_service = init_metrics_service(&data_dir).await?;
303
304        let startup_sessions = {
305            let entries = session_store.list_index_entries().await;
306            let mut sessions = Vec::new();
307            for entry in entries {
308                if let Some(session) = session_store
309                    .load_session(&entry.id)
310                    .await
311                    .map_err(AppError::StorageError)?
312                {
313                    sessions.push(session);
314                }
315            }
316            sessions
317        };
318        metrics_service
319            .reconcile_startup_sessions(startup_sessions, &[])
320            .await
321            .map_err(|error| {
322                AppError::InternalError(anyhow::anyhow!(
323                    "Failed to reconcile stale metrics state on startup: {error}"
324                ))
325            })?;
326
327        let agent_runners: Arc<RwLock<HashMap<String, AgentRunner>>> =
328            Arc::new(RwLock::new(HashMap::new()));
329        // NOTE: the idle-eviction sweep for `agent_runners` is spawned below,
330        // once `session_event_senders` also exists, so it can drop both maps'
331        // entries for a completed session together (issue #346).
332
333        let process_registry = Arc::new(ProcessRegistry::new());
334        let (provider_lock, provider_handle) = build_provider_handles(provider);
335
336        // Initialize multi-provider registry (for features.provider_model_ref).
337        let config_snapshot = config.read().await;
338        let provider_registry = match bamboo_llm::ProviderRegistry::from_config(
339            &config_snapshot,
340            bamboo_home_dir.clone(),
341        )
342        .await
343        {
344            Ok(registry) => Arc::new(registry),
345            Err(e) => {
346                tracing::error!("Failed to create provider registry: {}", e);
347                Arc::new(
348                    bamboo_llm::ProviderRegistry::from_config(
349                        &Config::default(),
350                        bamboo_home_dir.clone(),
351                    )
352                    .await
353                    .expect("Cannot create even an empty provider registry"),
354                )
355            }
356        };
357        drop(config_snapshot);
358
359        let provider_router = Arc::new(bamboo_llm::ProviderModelRouter::new(
360            provider_registry.clone(),
361        ));
362        let model_catalog = Arc::new(bamboo_llm::ModelCatalogService::new(
363            provider_registry.clone(),
364        ));
365
366        // Long-lived session event senders map (UI subscriptions + background tasks).
367        // Declared before `build_base_tools` (moved up from its original spot below)
368        // because the `notify` tool overlaid there needs it to broadcast onto a
369        // session's live channel — see `app_state::tools::build_base_tools`.
370        let session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>> =
371            Arc::new(RwLock::new(HashMap::new()));
372
373        // Shared bundle of always-on notification relay deps (see
374        // `session_events::NotificationRelayDeps`). Built once and cloned into
375        // every entry point that starts a relay directly at execution time —
376        // the schedule manager, the root child-session adapter, and the
377        // guardian child-session adapter below — so they can never drift.
378        let notification_relay_deps = crate::app_state::session_events::NotificationRelayDeps {
379            notification_service: notification_service.clone(),
380            session_event_senders: session_event_senders.clone(),
381            session_watchers: session_watchers.clone(),
382            config: config.clone(),
383        };
384
385        // The `ledger` tool needs the schedule store (built further down) to
386        // sync reminders; hand it a late-bound bridge now and bind it below.
387        let ledger_schedule_bridge =
388            Arc::new(crate::schedule_app::LateBoundLedgerBridge::default());
389
390        // The runtime and skill tools must share one cache-aware coordinator.
391        // Session setup publishes this run's resolved skill allowlist through it
392        // before the model can call load_skill.
393        let session_repo = bamboo_engine::SessionRepository::new(
394            sessions.clone(),
395            storage.clone(),
396            persistence.clone(),
397        );
398
399        // Account-scoped durable change feed. It is initialized before the
400        // Project tool surface so non-HTTP Project mutations publish the same
401        // replayable events as the REST API.
402        let account_sink = bamboo_engine::events::AccountEventSink::new(data_dir.join("events"))
403            .map_err(|e| {
404                AppError::InternalError(anyhow::anyhow!(
405                    "failed to initialize account change-feed journal: {e}"
406                ))
407            })?;
408        let project_resource_watcher = super::project_watcher::ProjectResourceWatcher::start(
409            project_store.clone(),
410            account_sink.clone(),
411            std::time::Duration::from_millis(120),
412        )
413        .map_err(|error| {
414            AppError::InternalError(anyhow::anyhow!(
415                "failed to start Project resource watcher: {error}"
416            ))
417        })?;
418
419        let base_tools = build_base_tools(
420            config.clone(),
421            permission_checker.clone(),
422            mcp_manager.clone(),
423            skill_manager.clone(),
424            session_repo.clone(),
425            bamboo_home_dir.clone(),
426            notification_service.clone(),
427            session_event_senders.clone(),
428            session_watchers.clone(),
429            ledger_schedule_bridge.clone(),
430            project_store.clone(),
431            account_sink.clone(),
432            workspace_resolver.clone(),
433        );
434
435        // The workflow engine executes against the base tool surface. The
436        // caller-facing workflow_run tool is overlaid onto the root surface
437        // later, preventing a workflow from recursively dispatching itself.
438        let workflow_runs = crate::workflow::WorkflowRunAccess::new_with_permission_config(
439            &data_dir,
440            base_tools.clone(),
441            skill_manager.clone(),
442            session_repo.clone(),
443            permission_checker.permission_config(),
444        )
445        .await
446        .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?;
447
448        // Idle-evict completed runners together with their paired session event
449        // senders (issue #346). Spawned here (not next to `agent_runners`) so it
450        // owns handles to both maps.
451        spawn_session_map_cleanup_task(agent_runners.clone(), session_event_senders.clone(), None);
452
453        // Bridge both instruction and orchestration catalog transitions onto
454        // the same durable account feed used by SSE and v2 WebSocket clients.
455        // The public Workflow Library unifies both kinds, so either namespace
456        // must invalidate clients. Catalog events are account-scoped (no
457        // session id).
458        {
459            let mut workflow_events = skill_manager.store().subscribe_workflow_catalog();
460            let account_sink = account_sink.clone();
461            tokio::spawn(async move {
462                loop {
463                    let event = match workflow_events.recv().await {
464                        Ok(event) => event,
465                        Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
466                            tracing::warn!("Workflow catalog event bridge lagged by {skipped}");
467                            continue;
468                        }
469                        Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
470                    };
471                    let event = match event.kind {
472                        bamboo_skills::WorkflowCatalogEventKind::Changed => {
473                            AgentEvent::WorkflowChanged {
474                                workflow_id: event.workflow_id,
475                                revision: event.revision,
476                                scope: event.scope,
477                            }
478                        }
479                        bamboo_skills::WorkflowCatalogEventKind::Invalid => {
480                            AgentEvent::WorkflowInvalid {
481                                workflow_id: event.workflow_id,
482                                revision: event.revision,
483                                scope: event.scope,
484                            }
485                        }
486                        bamboo_skills::WorkflowCatalogEventKind::Recovered => {
487                            AgentEvent::WorkflowRecovered {
488                                workflow_id: event.workflow_id,
489                                revision: event.revision,
490                                scope: event.scope,
491                            }
492                        }
493                    };
494                    account_sink.record(None, &event);
495                }
496            });
497        }
498        let (approval_registry, restart_approval_events) =
499            bamboo_engine::external_agents::live::initialize_durable_approvals(
500                data_dir.join("approvals/child-approvals-v1.json"),
501            )
502            .map_err(|error| {
503                AppError::InternalError(anyhow::anyhow!(
504                    "failed to initialize durable child approvals: {error}"
505                ))
506            })?;
507        for event in restart_approval_events {
508            account_sink.record(event.session_id(), &event);
509        }
510
511        // Sub-agents are full agents with the full toolset (no per-role tool
512        // trimming): the child tool surface is the plain base tools.
513        let child_tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = base_tools.clone();
514
515        // Unified agent runtime (shared resources for all execution paths).
516        // default_tools = base_tools (builtin + MCP + memory + skills) as a safe fallback.
517        // Interactive execution paths pass an explicit tool surface override:
518        // root sessions use ToolSurface::Root; child sessions use ToolSurface::Child.
519        let project_context_resolver = Arc::new(
520            bamboo_engine::project_context::ProjectContextResolver::new_with_workspace_resolver(
521                Arc::new(crate::project_context::ProjectStoreContextSource::new(
522                    project_store.clone(),
523                )),
524                workspace_resolver.clone(),
525            ),
526        );
527        let mut agent_builder = bamboo_engine::Agent::builder()
528            .storage(storage.clone())
529            .persistence(Arc::new(session_repo.clone()))
530            .session_inbox(session_inbox.clone())
531            .activation_router(session_activation_router.clone())
532            .session_messenger(session_messenger.clone())
533            .attachment_reader(session_store.clone())
534            .skill_manager(skill_manager.clone())
535            .metrics_collector(metrics_service.collector())
536            .config(config.clone())
537            .provider(provider_handle.clone())
538            .default_tools(base_tools.clone())
539            .project_context_resolver(project_context_resolver.clone());
540        if let Some(permission_config) = permission_checker.permission_config() {
541            agent_builder = agent_builder.permission_config(permission_config);
542        }
543        let agent = Arc::new(
544            agent_builder
545                .build()
546                .expect("agent runtime should be fully configured"),
547        );
548
549        let child_completion_coordinator =
550            Arc::new(bamboo_engine::ChildCompletionCoordinator::new(
551                storage.clone(),
552                persistence.clone(),
553                sessions.clone(),
554                agent_runners.clone(),
555                session_event_senders.clone(),
556                agent.clone(),
557                config.clone(),
558                provider_registry.clone(),
559                provider_router.clone(),
560                data_dir.clone(),
561                Some(account_sink.inbox()),
562            ));
563        session_activation_router
564            .set_spawner(child_completion_coordinator.clone())
565            .await;
566
567        // Initialize sub-session spawn scheduler (async background jobs).
568        let config_snapshot = config.read().await.clone();
569
570        // When a broker is configured, run the MCP proxy service under a
571        // supervisor (issue #47): deployed workers forward their (host-bound)
572        // MCP tool calls here, and we execute them against this orchestrator's
573        // real MCP servers (single MCP host). The supervisor restarts the proxy
574        // with bounded backoff after a transient WebSocket drop instead of
575        // permanently disabling proxied tools for the worker's lifetime.
576        let mcp_proxy_shutdown = tokio_util::sync::CancellationToken::new();
577        if let Some(broker) = config_snapshot.subagents().broker.clone() {
578            if !broker.endpoint.trim().is_empty() {
579                let backend: std::sync::Arc<dyn bamboo_agent_core::tools::ToolExecutor> =
580                    std::sync::Arc::new(bamboo_mcp::executor::McpToolExecutor::new(
581                        mcp_manager.clone(),
582                        mcp_manager.tool_index(),
583                    ));
584                let shutdown = mcp_proxy_shutdown.clone();
585
586                // Build the orchestrator-side per-role MCP tool allowlist from
587                // config (issue #54). Enforcement lives HERE — never in the
588                // worker-facing `McpProxyConfig` a deployed worker receives —
589                // because a worker self-declaring its own allowlist would be
590                // insecure (it could simply claim to be unrestricted). Validate
591                // configured tool names against THIS backend's real, live tool
592                // set so a typo in `config.json` is surfaced at boot instead of
593                // silently granting nothing for the intended tool.
594                let role_entries: Vec<(String, Vec<String>)> = config_snapshot
595                    .subagents()
596                    .mcp_role_allowlist
597                    .iter()
598                    .map(|e| (e.role.clone(), e.tools.clone()))
599                    .collect();
600                if role_entries.is_empty() {
601                    // Default behavior unchanged (issue #54 item 5): no policy
602                    // configured -> every role sees/can call every proxied
603                    // tool. Logged once here at boot (not per-request) so an
604                    // operator can tell the restriction is opt-in rather than
605                    // silently absent.
606                    tracing::info!(
607                        "mcp proxy: no subagents.mcp_role_allowlist configured — every worker \
608                         role sees/can call the full host-bound MCP tool set (opt in a role \
609                         policy in config.json to scope tools per role; see issue #54)"
610                    );
611                }
612                // NOTE: `init_mcp_manager` connects MCP servers in a background
613                // task so the HTTP API stays responsive at boot — this backend's
614                // `list_tools()` can legitimately be empty here if that task
615                // hasn't finished yet. `from_config` treats an empty set as "skip
616                // tool-name validation" rather than flagging every configured
617                // tool as an unknown typo, so warn separately here when that
618                // race means validation was effectively skipped.
619                let known_tools: std::collections::HashSet<String> = backend
620                    .list_tools()
621                    .into_iter()
622                    .map(|t| t.function.name)
623                    .collect();
624                if !role_entries.is_empty() && known_tools.is_empty() {
625                    tracing::warn!(
626                        "mcp role allowlist: the orchestrator's MCP tool set was empty at \
627                         policy-load time (servers may still be connecting in the background) — \
628                         skipped tool-name typo validation for subagents.mcp_role_allowlist"
629                    );
630                }
631                let allowlist = std::sync::Arc::new(bamboo_broker::RoleToolAllowlist::from_config(
632                    role_entries,
633                    &known_tools,
634                ));
635                tokio::spawn(async move {
636                    let me = bamboo_broker::AgentRef {
637                        session_id: bamboo_broker::ORCHESTRATOR_ID.to_string(),
638                        role: Some("orchestrator".into()),
639                    };
640                    bamboo_broker::serve_mcp_proxy_supervised(
641                        &broker.endpoint,
642                        me,
643                        &broker.token,
644                        backend,
645                        allowlist,
646                        shutdown,
647                    )
648                    .await;
649                });
650            }
651        }
652        let parent_approval_reviewer = Arc::new(
653            crate::app_state::parent_approval_reviewer::ParentAgentApprovalReviewer::new(
654                session_repo.clone(),
655                provider_router.clone(),
656            ),
657        );
658        let codex_run_tokens = Arc::new(crate::codex_run_tokens::CodexRunTokenRegistry::default());
659        let external_runner =
660            bamboo_engine::external_agents::runtime::build_external_child_runner_with_codex_tokens(
661                &config_snapshot,
662                Some(approval_registry.clone()),
663                Some(parent_approval_reviewer),
664                permission_checker.permission_config(),
665                Some(codex_run_tokens.clone()),
666            );
667        external_runner.set_session_inbox_runtime(Some(
668            bamboo_engine::execution::spawn::SessionInboxRuntimeBinding {
669                router: session_activation_router.clone(),
670                inbox: session_inbox.clone(),
671                storage: storage.clone(),
672                persistence: persistence.clone(),
673            },
674        ));
675        let spawn_scheduler = build_spawn_scheduler(
676            agent.clone(),
677            child_tools,
678            sessions.clone(),
679            agent_runners.clone(),
680            session_event_senders.clone(),
681            external_runner,
682            Some(provider_router.clone()),
683            Some(child_completion_coordinator.clone()),
684            Some(data_dir.clone()),
685            Some(account_sink.inbox()),
686            Some(Arc::new(
687                crate::app_state::session_events::NotificationRelayLaunchHook::new(
688                    notification_relay_deps.clone(),
689                ),
690            )),
691        );
692        child_completion_coordinator
693            .set_spawn_scheduler(&spawn_scheduler)
694            .await;
695
696        let tools_with_task = base_tools.clone();
697
698        let schedule_store = init_schedule_store(&data_dir).await?;
699
700        // Bind the ledger's reminder bridge now that the schedule store exists.
701        ledger_schedule_bridge
702            .bind(Arc::new(crate::schedule_app::ScheduleLedgerBridge::new(
703                schedule_store.clone(),
704            )))
705            .await;
706
707        let schedule_manager = build_schedule_manager(
708            schedule_store.clone(),
709            agent.clone(),
710            tools_with_task.clone(),
711            permission_checker.permission_config(),
712            sessions.clone(),
713            agent_runners.clone(),
714            session_event_senders.clone(),
715            persistence.clone(),
716            config.clone(),
717            provider_registry.clone(),
718            Some(data_dir.clone()),
719            Some(account_sink.inbox()),
720            notification_relay_deps.clone(),
721            project_store.clone(),
722            workspace_resolver.clone(),
723        );
724
725        bamboo_engine::auto_dream::spawn_auto_dream_task_with_project_resolver(
726            bamboo_engine::auto_dream::AutoDreamContext {
727                session_store: session_store.clone(),
728                storage: storage.clone(),
729                provider: provider_handle.clone(),
730                config: config.clone(),
731                provider_registry: provider_registry.clone(),
732            },
733            project_context_resolver.as_ref().clone(),
734        );
735
736        // Background memory "gardener": opt-in blob remediation + near-duplicate
737        // consolidation. No-op cost unless `memory.gardener_enabled` /
738        // `memory.dedup_gardener_enabled` is set; an empty prefilter makes zero LLM calls.
739        bamboo_engine::gardener::spawn_gardener_task_with_project_resolver(
740            bamboo_engine::auto_dream::AutoDreamContext {
741                session_store: session_store.clone(),
742                storage: storage.clone(),
743                provider: provider_handle.clone(),
744                config: config.clone(),
745                provider_registry: provider_registry.clone(),
746            },
747            project_context_resolver.clone(),
748        );
749
750        // Background ledger gardener: expiry + record↔schedule reconciliation
751        // are deterministic and free; distillation uses the background model
752        // and no-ops without one. The bridge handle is already bound above.
753        bamboo_engine::ledger_gardener::spawn_ledger_gardener_task(
754            bamboo_engine::ledger_gardener::LedgerGardenerContext {
755                dream: bamboo_engine::auto_dream::AutoDreamContext {
756                    session_store: session_store.clone(),
757                    storage: storage.clone(),
758                    provider: provider_handle.clone(),
759                    config: config.clone(),
760                    provider_registry: provider_registry.clone(),
761                },
762                schedule_bridge: Some(ledger_schedule_bridge.clone()),
763            },
764        );
765
766        let config_for_resolver = config.clone();
767        let subagent_model_resolver: OptionalSubagentModelResolver = {
768            let registry = provider_registry.clone();
769            Some(Arc::new(
770                move |subagent_type: String| -> futures::future::BoxFuture<
771                    'static,
772                    Option<bamboo_domain::ProviderModelRef>,
773                > {
774                    let config_for_resolver = config_for_resolver.clone();
775                    let registry = registry.clone();
776                    Box::pin(async move {
777                        let config_snap = config_for_resolver.read().await.clone();
778                        bamboo_engine::model_config_helper::resolve_subagent_model_ref(
779                            &config_snap,
780                            &config_snap.provider,
781                            &registry,
782                            &subagent_type,
783                        )
784                    })
785                },
786            ))
787        };
788
789        // Config-write io-lock + the shared Remote Cluster Fabric deploy engine.
790        // The engine is built once and shared by the HTTP handlers (via AppState)
791        // and the `cluster` agent tool, so both use ONE worker registry.
792        let config_io_lock = Arc::new(tokio::sync::Mutex::new(()));
793        let fabric_registry: crate::tools::DeployedRegistry =
794            Arc::new(tokio::sync::Mutex::new(HashMap::new()));
795        let fabric_bamboo_bin =
796            std::env::current_exe().unwrap_or_else(|_| std::path::PathBuf::from("bamboo"));
797        let credential_store = Arc::new(bamboo_config::CredentialStore::open(&bamboo_home_dir));
798        let mut fabric_deployer = bamboo_server_tools::FabricDeployer::new(
799            config.clone(),
800            config_io_lock.clone(),
801            bamboo_home_dir.clone(),
802            fabric_registry,
803            fabric_bamboo_bin,
804        );
805        if let Some(facade) = config_facade.clone() {
806            let event_sink = account_sink.clone();
807            fabric_deployer = fabric_deployer.with_modular_persistence(
808                facade,
809                credential_store.clone(),
810                Arc::new(move |event| {
811                    let event_sink = event_sink.clone();
812                    let event = event.clone();
813                    Box::pin(async move {
814                        super::config_runtime::publish_registry_event(&event_sink, &event).await;
815                    })
816                }),
817            );
818        }
819        if let Err(error) = fabric_deployer.reconcile_stale_nodes_on_boot().await {
820            tracing::warn!(
821                error = %error,
822                "cluster-fabric boot reconcile failed; retaining the last adopted state"
823            );
824        }
825        let fabric_deployer = Arc::new(fabric_deployer);
826        // Cluster health monitor: periodically probe deployed workers on the bus and
827        // flip node status live (Running↔Unreachable) + auto-recover. Server-scoped
828        // — it runs under BOTH the embedded and an external broker (it reads the
829        // broker endpoint lazily each tick), and is aborted when the server drops.
830        let health_monitor = fabric_deployer
831            .clone()
832            .spawn_health_monitor()
833            .await
834            .map(HealthMonitor);
835
836        let tools = build_root_tools(
837            tools_with_task.clone(),
838            schedule_store.clone(),
839            schedule_manager.clone(),
840            session_store.clone(),
841            storage.clone(),
842            persistence.clone(),
843            session_messenger.clone(),
844            spawn_scheduler.clone(),
845            sessions.clone(),
846            agent_runners.clone(),
847            session_event_senders.clone(),
848            subagent_model_resolver,
849            config.clone(),
850            provider_registry.clone(),
851            config_snapshot.subagents().broker.clone(),
852            fabric_deployer.clone(),
853            project_store.clone(),
854            workspace_resolver.clone(),
855        );
856        let workflow_run_tool =
857            Arc::new(crate::workflow::WorkflowRunTool::new(workflow_runs.clone()));
858        let tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor> = Arc::new(
859            crate::tools::OverlayToolExecutor::new(tools, workflow_run_tool),
860        );
861
862        child_completion_coordinator
863            .set_root_tools(tools.clone())
864            .await;
865
866        // Process restart recovery: only backlog covered by its producer's
867        // durable activation watermark requests a run. A child/Bash coordinator
868        // may intentionally stage sibling outcomes while a specific wait remains
869        // armed; admission by itself is not permission to execute.
870        for entry in session_store.list_index_entries().await {
871            match session_inbox.inspect(&entry.id).await {
872                Ok(backlog) if backlog.activation_pending() => {
873                    if let Err(error) = bamboo_domain::SessionActivationPort::request_activation(
874                        session_activation_router.as_ref(),
875                        &entry.id,
876                        backlog.activation_generation,
877                    )
878                    .await
879                    {
880                        tracing::error!(
881                            session_id = %entry.id,
882                            %error,
883                            "failed to reactivate durable SessionInbox backlog during startup"
884                        );
885                    }
886                }
887                Ok(_) => {}
888                Err(error) => tracing::warn!(
889                    session_id = %entry.id,
890                    %error,
891                    "failed to inspect SessionInbox during startup recovery"
892                ),
893            }
894        }
895
896        let tool_factory =
897            crate::tools::ToolSurfaceFactory::new(base_tools, tools_with_task, tools);
898
899        let session_repo = bamboo_engine::SessionRepository::new(
900            sessions.clone(),
901            storage.clone(),
902            persistence.clone(),
903        );
904
905        // bamboo-connect (#452 / epic #447): drives bamboo sessions from IM
906        // platforms (Telegram first). Fully inert when `config.connect.platforms`
907        // is empty — mirrors the schedule manager / notification relay's
908        // always-constructed-but-only-active-if-configured lifecycle.
909        let connect_manager = Arc::new(
910            build_connect_manager(
911                agent.clone(),
912                tool_factory.get(crate::tools::ToolSurface::Root),
913                session_repo.clone(),
914                agent_runners.clone(),
915                session_event_senders.clone(),
916                Some(account_sink.inbox()),
917                Some(data_dir.clone()),
918                config.clone(),
919                provider_registry.clone(),
920                permission_checker.clone(),
921                project_store.clone(),
922                workspace_resolver.clone(),
923            )
924            .await
925            .map_err(|error| AppError::InternalError(anyhow::anyhow!(error)))?,
926        );
927
928        // Dedicated child-session adapter backing the guardian review spawner.
929        // The guardian path passes an explicit model (no subagent_type routing)
930        // and registers its parent wait at the terminal gate (not via the
931        // adapter's coalescing slots), so a lightweight adapter with no resolver
932        // and a fresh wait-slot map suffices. `Arc<ChildSessionAdapter>` doubles
933        // as `Arc<dyn GuardianSpawner>`.
934        let child_adapter = Arc::new(crate::tools::ChildSessionAdapter {
935            session_store: session_store.clone(),
936            storage: storage.clone(),
937            persistence: persistence.clone(),
938            session_messenger: Some(session_messenger.clone()),
939            scheduler: spawn_scheduler.clone(),
940            sessions_cache: sessions.clone(),
941            agent_runners: agent_runners.clone(),
942            session_event_senders: session_event_senders.clone(),
943            subagent_model_resolver: None,
944            config: config.clone(),
945            project_store: Some(project_store.clone()),
946            workspace_resolver: workspace_resolver.clone(),
947            parent_wait_slots: Arc::new(dashmap::DashMap::new()),
948        });
949        let guardian_spawner: Arc<dyn bamboo_engine::GuardianSpawner> = child_adapter.clone();
950        // Wire the spawner into the completion coordinator too, so a resumed run
951        // can re-spawn a guardian to re-review a fix after a reject verdict.
952        child_completion_coordinator
953            .set_guardian_spawner(guardian_spawner.clone())
954            .await;
955
956        // The completion coordinator doubles as the bash self-resume hook
957        // (issue #84 Phase 2b): it polls the live shell registry and resumes a
958        // session once all its background bash shells finish.
959        let bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook> =
960            child_completion_coordinator.clone();
961
962        // Child-wait watchdog (issue #546): boot-time reconciliation of
963        // children orphaned by a restart, then a periodic heartbeat sweep that
964        // backstops every lost child→parent wake (panicked child task, dead
965        // spawn scheduler, clobbered resume, expired wait lease, ...). Spawned
966        // AFTER `set_root_tools` above so a boot-time parent resume can spawn.
967        child_completion_coordinator.spawn_child_wait_watchdog();
968
969        // `ServiceManager` (issue #479 / epic #477 prereq): supervises
970        // long-running "service" plugins. Always constructed — fully inert
971        // until a plugin install or the boot-time reconcile below starts
972        // something — mirrors `mcp_manager`/`connect_manager`'s
973        // always-alive lifecycle.
974        let service_manager = Arc::new(crate::service_manager::ServiceManager::new());
975        // Backgrounded (mirrors `init_mcp_manager`'s background MCP
976        // bootstrap): a service that `installed.json` says should be
977        // running but isn't (the previous `bamboo serve` process, if
978        // any, died with everything it supervised) is started fresh.
979        //
980        // The `JoinHandle` is kept (not discarded) purely so tests can
981        // deterministically wait it out via
982        // `AppState::wait_for_boot_reconcile_services` instead of racing
983        // this unsynchronized pass — see that method's doc comment and issue
984        // #486. Production code never awaits it; server startup is never
985        // blocked on plugin service spawns.
986        let boot_reconcile_services_handle = {
987            let service_manager = service_manager.clone();
988            let app_data_dir = bamboo_home_dir.clone();
989            tokio::spawn(async move {
990                crate::plugin_installer::boot_reconcile_services(&app_data_dir, &service_manager)
991                    .await;
992            })
993        };
994
995        let (config_watcher, config_live_health, mcp_config_live_health) =
996            super::config_runtime::ConfigWatcherRuntime::start(
997                bamboo_home_dir.clone(),
998                config.clone(),
999                config_facade.clone(),
1000                config_io_lock.clone(),
1001                provider_registry.clone(),
1002                provider_lock.clone(),
1003                mcp_manager.clone(),
1004                account_sink.clone(),
1005            );
1006        Ok(Self {
1007            app_data_dir: bamboo_home_dir,
1008            config,
1009            config_facade,
1010            config_io_lock,
1011            config_live_health,
1012            mcp_config_live_health,
1013            config_watcher,
1014            project_resource_watcher,
1015            credential_store,
1016            fabric_deployer,
1017            embedded_broker,
1018            health_monitor,
1019            provider: provider_lock,
1020            provider_handle,
1021            sessions,
1022            storage,
1023            session_store,
1024            session_create_operations,
1025            mutation_idempotency,
1026            project_store,
1027            project_context_resolver,
1028            workspace_resolver,
1029            session_repo,
1030            persistence,
1031            session_inbox,
1032            session_activation_router,
1033            session_messenger,
1034            spawn_scheduler,
1035            child_completion_coordinator,
1036            guardian_spawner,
1037            bash_resume_hook,
1038            schedule_store,
1039            schedule_manager,
1040            connect_manager,
1041            tool_factory,
1042            permission_checker,
1043            permission_section,
1044            permission_io_lock,
1045            approval_registry,
1046            notification_service,
1047            session_watchers,
1048            cancel_tokens: Arc::new(RwLock::new(HashMap::new())),
1049            mcp_proxy_shutdown,
1050            skill_manager,
1051            workflow_runs,
1052            mcp_manager,
1053            service_manager,
1054            boot_reconcile_services_handle: tokio::sync::Mutex::new(Some(
1055                boot_reconcile_services_handle,
1056            )),
1057            metrics_service,
1058            agent_runners,
1059            execute_startups: Arc::new(std::sync::Mutex::new(HashMap::new())),
1060            session_event_senders,
1061            account_sink,
1062            process_registry,
1063            metrics_bus: None, // Will be set by server if needed
1064            agent,
1065            provider_registry,
1066            provider_router,
1067            model_catalog,
1068            title_gen_in_flight: Arc::new(dashmap::DashSet::new()),
1069            pairing_codes: Arc::new(dashmap::DashMap::new()),
1070            pairing_code_guard: Arc::new(crate::handlers::settings::PairingCodeGuard::default()),
1071            root_password_guard: Arc::new(crate::handlers::settings::RootPasswordGuard::default()),
1072            codex_run_tokens,
1073            // remote-actor P2a (#181): empty in-memory agent registry.
1074        })
1075    }
1076}
1077
1078/// A handle to the in-process mailbox bus (broker) so it can be shut down with
1079/// the server. Dropping it aborts the serve task.
1080pub struct EmbeddedBroker {
1081    task: tokio::task::JoinHandle<()>,
1082    gc_task: tokio::task::JoinHandle<()>,
1083}
1084
1085impl Drop for EmbeddedBroker {
1086    fn drop(&mut self) {
1087        self.task.abort();
1088        self.gc_task.abort();
1089    }
1090}
1091
1092/// Server-scoped handle to the cluster health monitor. Kept separate from
1093/// [`EmbeddedBroker`] so the monitor runs under BOTH the embedded broker and an
1094/// external (`broker.json`) one; aborts the sweep on drop.
1095pub struct HealthMonitor(tokio::task::JoinHandle<()>);
1096
1097impl Drop for HealthMonitor {
1098    fn drop(&mut self) {
1099        self.0.abort();
1100    }
1101}
1102
1103/// Start an in-process broker on `127.0.0.1:<auto>` and point
1104/// `config.subagents.broker` (RUNTIME-ONLY, `#[serde(skip)]`) at it — UNLESS a
1105/// user-managed external broker is configured in `<data_dir>/broker.json` and
1106/// reachable (then use that). Returns `None` when an external broker is used or
1107/// the bind fails (sub-agent dispatch then degrades exactly as before).
1108///
1109/// The broker's endpoint deliberately lives in EITHER the in-memory config (for
1110/// the embedded case, regenerated each boot) or its own `broker.json` (for the
1111/// external case) — NEVER in `config.json`. That is what stops a prior run's
1112/// ephemeral auto-port from leaking into the user's config and being dialed dead
1113/// on the next boot (every sub-agent + the MCP proxy would hit "connect refused").
1114async fn maybe_embed_broker(
1115    config: &mut bamboo_llm::Config,
1116    data_dir: &std::path::Path,
1117) -> Option<EmbeddedBroker> {
1118    // A user-managed EXTERNAL broker (multi-host / shared standalone bus) lives in
1119    // its OWN file, `<data_dir>/broker.json` — separate from config.json. Honour
1120    // it ONLY if actually REACHABLE; a dead endpoint (standalone broker not up)
1121    // falls through to a fresh in-process broker so dispatch still works.
1122    if let Some(external) = load_external_broker(data_dir) {
1123        let endpoint = external.endpoint.trim().to_string();
1124        if broker_endpoint_reachable(&endpoint).await {
1125            tracing::info!(%endpoint, "using external broker from broker.json");
1126            config.subagents_mut().broker = Some(external);
1127            return None;
1128        }
1129        tracing::warn!(
1130            %endpoint,
1131            "broker.json endpoint is unreachable — embedding a fresh in-process broker instead"
1132        );
1133    }
1134
1135    let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
1136        Ok(l) => l,
1137        Err(e) => {
1138            tracing::warn!("embedded broker: bind failed, sub-agent dispatch disabled: {e}");
1139            return None;
1140        }
1141    };
1142    let port = match listener.local_addr() {
1143        Ok(a) => a.port(),
1144        Err(e) => {
1145            tracing::warn!("embedded broker: local_addr failed: {e}");
1146            return None;
1147        }
1148    };
1149    let token = uuid::Uuid::new_v4().simple().to_string();
1150    let root = data_dir.join("broker");
1151    let core = Arc::new(bamboo_broker::BrokerCore::new(root));
1152    // Reclaim orphan mailbox dirs (one-shot parent links, killed pool workers)
1153    // every 5 min so `<data>/broker/mailboxes/` doesn't grow unbounded.
1154    let gc_task = core
1155        .clone()
1156        .spawn_mailbox_gc(std::time::Duration::from_secs(300));
1157    let server = Arc::new(bamboo_broker::BrokerServer::new(core, token.clone()));
1158
1159    let task = tokio::spawn(async move {
1160        if let Err(e) = server.serve(listener).await {
1161            tracing::error!("embedded broker serve loop ended: {e}");
1162        }
1163    });
1164
1165    // Set the endpoint IN MEMORY ONLY — `subagents.broker` is `#[serde(skip)]`, so
1166    // this ephemeral loopback port is regenerated every boot and never touches disk.
1167    config.subagents_mut().broker = Some(bamboo_config::BrokerClientConfig {
1168        endpoint: format!("ws://127.0.0.1:{port}"),
1169        token,
1170        token_encrypted: None,
1171        credential_ref: None,
1172        configured: false,
1173    });
1174    tracing::info!(port, "embedded mailbox bus (broker) started in-process");
1175    Some(EmbeddedBroker { task, gc_task })
1176}
1177
1178/// Load a user-managed EXTERNAL broker from `<data_dir>/broker.json`, if present.
1179/// This file is the SEPARATE, persisted home for a standalone/remote broker —
1180/// deliberately NOT `config.json`, so the embedded broker's ephemeral runtime
1181/// port can never leak into the user's config (the stale-dead-port bug). An
1182/// absent file or a parse error yields `None` (embed a fresh in-process broker).
1183///
1184/// Durable format carries only endpoint plus credential metadata. Legacy
1185/// plaintext/ciphertext tokens are migrated before this reader proceeds.
1186fn load_external_broker(data_dir: &std::path::Path) -> Option<bamboo_config::BrokerClientConfig> {
1187    let path = data_dir.join("broker.json");
1188    if let Err(error) = bamboo_config::migrate_external_broker_credentials(data_dir)
1189        .and_then(|_| bamboo_config::ensure_provider_mcp_migration_ready(data_dir))
1190    {
1191        tracing::warn!(error = %error, "external broker credential migration unavailable");
1192        return None;
1193    }
1194    let bytes = std::fs::read(&path).ok()?;
1195    match parse_external_broker_snapshot(&bytes, data_dir) {
1196        Ok(cfg) => Some(cfg),
1197        Err(error) => {
1198            tracing::warn!(?path, %error, "broker.json is unavailable");
1199            None
1200        }
1201    }
1202}
1203
1204fn parse_external_broker_snapshot(
1205    bytes: &[u8],
1206    data_dir: &std::path::Path,
1207) -> bamboo_config::ConfigStoreResult<bamboo_config::BrokerClientConfig> {
1208    let mut cfg: bamboo_config::BrokerClientConfig = serde_json::from_slice(bytes)?;
1209    if !cfg.token.trim().is_empty()
1210        || cfg
1211            .token_encrypted
1212            .as_deref()
1213            .is_some_and(|value| !value.trim().is_empty())
1214    {
1215        return Err(bamboo_config::ConfigStoreError::Validation(
1216            "legacy broker credential appeared after migration".to_string(),
1217        ));
1218    }
1219    if cfg.endpoint.trim().is_empty() {
1220        return Err(bamboo_config::ConfigStoreError::Validation(
1221            "broker endpoint is empty".to_string(),
1222        ));
1223    }
1224    cfg.hydrate_credential_from_store(data_dir)?;
1225    Ok(cfg)
1226}
1227
1228/// Best-effort TCP reachability probe of a `ws[s]://host:port[/path]` broker
1229/// endpoint. Used to tell a LIVE external/standalone broker (keep) apart from a
1230/// DEAD persisted endpoint (a prior run's embedded auto-port that leaked into
1231/// config.json — re-embed). A short timeout keeps boot fast when it's dead.
1232async fn broker_endpoint_reachable(endpoint: &str) -> bool {
1233    let host_port = endpoint
1234        .trim()
1235        .trim_start_matches("wss://")
1236        .trim_start_matches("ws://")
1237        .split('/')
1238        .next()
1239        .unwrap_or("");
1240    if host_port.is_empty() {
1241        return false;
1242    }
1243    matches!(
1244        tokio::time::timeout(
1245            std::time::Duration::from_millis(500),
1246            tokio::net::TcpStream::connect(host_port),
1247        )
1248        .await,
1249        Ok(Ok(_))
1250    )
1251}
1252
1253#[cfg(test)]
1254mod fabric_boot_reconcile_tests {
1255    use super::*;
1256    use bamboo_config::cluster_fabric::{
1257        DeployProfile, Node, NodePlacement, NodeState, NodeStatus, TrustLevel,
1258    };
1259    use bamboo_config::{ClusterNodeCredentialIntents, ConfigFacade};
1260    use std::collections::BTreeMap;
1261
1262    #[tokio::test]
1263    async fn restart_reconcile_keeps_runtime_process_facade_and_disk_on_one_revision() {
1264        let _key = bamboo_config::encryption::set_test_encryption_key([0x7b; 32]);
1265        let dir = tempfile::tempdir().unwrap();
1266        let first = AppState::new(dir.path().to_path_buf()).await.unwrap();
1267        first
1268            .update_cluster_fabric_credentials(
1269                0,
1270                BTreeMap::from([(
1271                    "boot-node".to_string(),
1272                    ClusterNodeCredentialIntents::clear_all(),
1273                )]),
1274                |config| {
1275                    config.cluster_fabric.nodes.push(Node {
1276                        id: "boot-node".to_string(),
1277                        label: "boot-node".to_string(),
1278                        placement: NodePlacement::Local,
1279                        trust_level: TrustLevel::Trusted,
1280                        deploy: DeployProfile::default(),
1281                        state: Some(NodeState {
1282                            status: NodeStatus::Running,
1283                            worker_id: Some("stale-worker".to_string()),
1284                            ..Default::default()
1285                        }),
1286                        enabled: true,
1287                    });
1288                    Ok(())
1289                },
1290            )
1291            .await
1292            .unwrap();
1293        assert_eq!(
1294            first
1295                .config_facade
1296                .as_ref()
1297                .unwrap()
1298                .registry()
1299                .cluster_fabric
1300                .snapshot()
1301                .revision,
1302            1
1303        );
1304        drop(first);
1305
1306        let restarted = AppState::new(dir.path().to_path_buf()).await.unwrap();
1307        let process_snapshot = restarted
1308            .config_facade
1309            .as_ref()
1310            .unwrap()
1311            .registry()
1312            .cluster_fabric
1313            .snapshot();
1314        assert_eq!(process_snapshot.revision, 2);
1315        assert_eq!(
1316            process_snapshot
1317                .data
1318                .0
1319                .node("boot-node")
1320                .unwrap()
1321                .state
1322                .as_ref()
1323                .unwrap()
1324                .status,
1325            NodeStatus::Unreachable
1326        );
1327        assert_eq!(
1328            restarted
1329                .config
1330                .read()
1331                .await
1332                .cluster_fabric
1333                .node("boot-node")
1334                .unwrap()
1335                .state
1336                .as_ref()
1337                .unwrap()
1338                .status,
1339            NodeStatus::Unreachable
1340        );
1341        let reopened = ConfigFacade::open(dir.path()).unwrap();
1342        assert_eq!(reopened.registry().cluster_fabric.snapshot().revision, 2);
1343        assert_eq!(
1344            reopened
1345                .effective_config()
1346                .cluster_fabric
1347                .node("boot-node")
1348                .unwrap()
1349                .state
1350                .as_ref()
1351                .unwrap()
1352                .status,
1353            NodeStatus::Unreachable
1354        );
1355    }
1356}
1357
1358#[cfg(test)]
1359mod broker_embed_tests {
1360    use super::{broker_endpoint_reachable, load_external_broker, parse_external_broker_snapshot};
1361
1362    #[test]
1363    fn load_external_broker_reads_broker_json_not_config() {
1364        let _key = bamboo_config::encryption::set_test_encryption_key([0x57; 32]);
1365        let dir = tempfile::tempdir().unwrap();
1366        // Absent file ⇒ None (embed a fresh in-process broker).
1367        assert!(load_external_broker(dir.path()).is_none());
1368
1369        // A well-formed broker.json ⇒ parsed external broker.
1370        std::fs::write(
1371            dir.path().join("broker.json"),
1372            r#"{ "endpoint": "wss://broker.example:9600", "token": "t" }"#,
1373        )
1374        .unwrap();
1375        let got = load_external_broker(dir.path()).expect("parsed");
1376        assert_eq!(got.endpoint, "wss://broker.example:9600");
1377        assert_eq!(got.token, "t");
1378        assert_eq!(
1379            got.credential_ref.as_ref().unwrap().as_str(),
1380            "broker.external.bearer_token"
1381        );
1382        assert!(got.configured);
1383        let durable = std::fs::read_to_string(dir.path().join("broker.json")).unwrap();
1384        assert!(!durable.contains("\"token\""));
1385        assert!(!durable.contains("token_encrypted"));
1386        assert!(!durable.contains("\"t\""));
1387
1388        // Empty endpoint ⇒ ignored (treated as absent).
1389        std::fs::write(dir.path().join("broker.json"), r#"{ "endpoint": "  " }"#).unwrap();
1390        assert!(load_external_broker(dir.path()).is_none());
1391
1392        // Garbage ⇒ ignored, never panics.
1393        std::fs::write(dir.path().join("broker.json"), "not json").unwrap();
1394        assert!(load_external_broker(dir.path()).is_none());
1395    }
1396
1397    #[test]
1398    fn external_broker_reference_fails_closed_and_tracks_generic_cas_updates() {
1399        let _key = bamboo_config::encryption::set_test_encryption_key([0x58; 32]);
1400        let dir = tempfile::tempdir().unwrap();
1401        let reference =
1402            bamboo_config::credential_ref("broker", "external", "bearer_token").unwrap();
1403        std::fs::write(
1404            dir.path().join("broker.json"),
1405            serde_json::to_vec_pretty(&serde_json::json!({
1406                "endpoint": "wss://broker.example:9600",
1407                "credential_ref": reference,
1408                "configured": true,
1409            }))
1410            .unwrap(),
1411        )
1412        .unwrap();
1413
1414        assert!(load_external_broker(dir.path()).is_none());
1415        std::fs::write(
1416            dir.path().join("broker.json"),
1417            serde_json::to_vec_pretty(&serde_json::json!({
1418                "endpoint": "wss://broker.example:9600",
1419                "credential_ref": reference,
1420                "configured": false,
1421            }))
1422            .unwrap(),
1423        )
1424        .unwrap();
1425        assert!(load_external_broker(dir.path()).is_none());
1426        let store = bamboo_config::CredentialStore::open(dir.path());
1427        store
1428            .replace(
1429                reference.clone(),
1430                "replacement-token",
1431                bamboo_config::CredentialSource::User,
1432                0,
1433            )
1434            .unwrap();
1435        let hydrated = load_external_broker(dir.path()).unwrap();
1436        assert_eq!(hydrated.token, "replacement-token");
1437        store.clear(&reference, 1).unwrap();
1438        assert!(load_external_broker(dir.path()).is_none());
1439    }
1440
1441    #[test]
1442    fn broker_snapshot_with_late_legacy_secret_fails_closed() {
1443        let _key = bamboo_config::encryption::set_test_encryption_key([0x59; 32]);
1444        let dir = tempfile::tempdir().unwrap();
1445        let error = parse_external_broker_snapshot(
1446            br#"{
1447                "endpoint": "wss://broker.example:9600",
1448                "token": "late-legacy-token"
1449            }"#,
1450            dir.path(),
1451        )
1452        .unwrap_err();
1453        assert!(error.to_string().contains("appeared after migration"));
1454    }
1455
1456    #[tokio::test]
1457    async fn reachability_probe_distinguishes_live_from_dead() {
1458        // Bound-then-dropped port: nothing listening ⇒ unreachable (a stale
1459        // persisted embedded endpoint must NOT be trusted).
1460        let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1461        let dead = l.local_addr().unwrap();
1462        drop(l);
1463        assert!(!broker_endpoint_reachable(&format!("ws://{dead}")).await);
1464
1465        // Empty / malformed ⇒ never trusted.
1466        assert!(!broker_endpoint_reachable("").await);
1467        assert!(!broker_endpoint_reachable("ws://").await);
1468
1469        // A LIVE listener (with a path) ⇒ reachable (a real external broker is kept).
1470        let live = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1471        let addr = live.local_addr().unwrap();
1472        assert!(broker_endpoint_reachable(&format!("ws://{addr}/stream")).await);
1473    }
1474}