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