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